mysqlsyncer 是一个基于 MySQL binlog 的实时数据同步工具,将 MySQL 数据库表的变更实时同步到 Elasticsearch 索引中。
通过伪装为 MySQL Slave 节点,监听 binlog 事件,自动完成全量数据 dump 和增量数据同步,支持断点续传和优雅关闭。
- 全量 + 增量同步:首次启动自动 dump 全量数据,之后持续监听 binlog 做增量同步
- 断点续传:同步进度持久化,重启后从上次位置继续,无需重新全量同步
- 多表映射:支持多张 MySQL 表到多个 ES 索引的灵活映射
- 可插拔存储:同步进度通过
PositionRepository接口持久化,内置 ES 实现,可扩展为 Redis、etcd 等后端 - 可扩展写入:通过
Writer接口抽象写入目标,内置 ES 实现,可扩展为 Kafka、MongoDB 等下游 - 字段过滤:支持白名单(FilterColumns)和黑名单(IgnoreColumns)两种模式
- 自定义处理器:每个映射可配置 Handler 函数,对行数据进行自定义过滤
- 批量写入:支持按数量(BatchNum)和时间间隔(SyncDuration)两种触发条件,先到先执行
- 优雅关闭:Close() 方法安全停止同步并持久化最后的进度
- 类型全覆盖:完整支持 MySQL 常用列类型(DATETIME、TIMESTAMP、DATE、TIME、ENUM、SET、BIT、JSON、BLOB、BINARY/VARBINARY 等)
| 依赖 | 说明 |
|---|---|
| Go | >= 1.25 |
| MySQL | >= 5.7,需开启 binlog(ROW 格式),帐号需 REPLICATION SLAVE、REPLICATION CLIENT 权限 |
| mysqldump | 必须安装,且版本需与 MySQL Server 版本一致(如 MySQL 8.0.43 应使用对应版本的 mysqldump,否则可能出现数据格式不兼容) |
| Elasticsearch | >= 8.0(使用 ES 客户端 v9) |
go get github.com/alberliu/mysqlsyncerNewESSyncer 是便捷构造函数,内部自动创建 position service 和 ES writer,适合大多数场景。
package main
import (
"fmt"
"log"
"os"
"github.com/elastic/go-elasticsearch/v9"
"github.com/alberliu/mysqlsyncer"
"github.com/alberliu/mysqlsyncer/writer"
)
func main() {
dsn := fmt.Sprintf("%s:%s@tcp(%s)/%s?parseTime=true&loc=Local",
os.Getenv("MYSQL_USER"),
os.Getenv("MYSQL_PASSWORD"),
os.Getenv("MYSQL_HOST"), // 如 "127.0.0.1:3306"
"mydb",
)
esClient, err := elasticsearch.NewTyped(
elasticsearch.WithAddresses(os.Getenv("ES_HOST")), // 如 "http://127.0.0.1:9200"
)
if err != nil {
log.Fatal(err)
}
syncer, err := mysqlsyncer.NewESSyncer(
dsn, // MySQL DSN
esClient, // ES 客户端
"mydb", // 数据库名(同时作为同步位点的命名空间)
[]writer.ESMapping{
{
TableName: "users",
Index: "users_index",
FilterColumns: []string{"id", "name", "email", "status"},
},
},
)
if err != nil {
log.Fatal(err)
}
// Run 会阻塞,在另一个 goroutine 中调用 Close() 可优雅退出
if err := syncer.Run(); err != nil {
log.Fatal(err)
}
}手动创建 position.Service 和 writer.ESWriter,再传给 NewSyncer。适合需要自定义 position 存储后端或 Writer 实现的高级场景。
package main
import (
"fmt"
"log"
"log/slog"
"os"
"time"
"github.com/elastic/go-elasticsearch/v9"
"github.com/alberliu/mysqlsyncer"
"github.com/alberliu/mysqlsyncer/position"
"github.com/alberliu/mysqlsyncer/writer"
)
func main() {
dsn := fmt.Sprintf("%s:%s@tcp(%s)/%s?parseTime=true&loc=Local",
os.Getenv("MYSQL_USER"),
os.Getenv("MYSQL_PASSWORD"),
os.Getenv("MYSQL_HOST"),
"mydb",
)
esClient, _ := elasticsearch.NewTyped(
elasticsearch.WithAddresses(os.Getenv("ES_HOST")),
)
// 1. 创建 position service(管理同步进度持久化)
ps, err := position.NewServiceWithES(esClient, "mydb", slog.Default())
if err != nil {
log.Fatal(err)
}
// 2. 创建 ES writer
esWriter, err := writer.NewESWriter(esClient, []writer.ESMapping{
{
TableName: "users",
Index: "users_index",
FilterColumns: []string{"id", "name", "email", "status"},
},
})
if err != nil {
log.Fatal(err)
}
// 3. 创建 syncer,传入可选配置
syncer, err := mysqlsyncer.NewSyncer(
dsn,
esWriter,
ps,
mysqlsyncer.WithBatchNum(200),
mysqlsyncer.WithSyncDuration(50*time.Millisecond),
mysqlsyncer.WithLogger(slog.Default()),
)
if err != nil {
log.Fatal(err)
}
if err := syncer.Run(); err != nil {
log.Fatal(err)
}
}| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
mysqlDSN |
string |
是 | MySQL 连接串,格式 user:password@tcp(host)/db?parseTime=true&loc=Local |
w |
writer.Writer |
是 | Writer 实例,负责将事件写入目标端 |
ps |
*position.Service |
是 | 同步进度管理服务 |
opts |
...Option |
否 | 可选配置项,见下方 |
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
mysqlDSN |
string |
是 | MySQL 连接串 |
esClient |
*elasticsearch.TypedClient |
是 | 已初始化的 ES v9 客户端 |
database |
string |
是 | 数据库名,同时作为同步位点的命名空间 |
mappings |
[]writer.ESMapping |
是 | 表到索引的映射配置 |
opts |
...Option |
否 | 可选配置项 |
| 函数 | 默认值 | 说明 |
|---|---|---|
WithBatchNum(n int) |
100 |
批量写入条数,累积到该数量后触发写入 |
WithSyncDuration(d time.Duration) |
100ms |
最大同步间隔,即使未达到 BatchNum,超过此时间也会触发写入 |
WithLogger(l *slog.Logger) |
slog.Default() |
结构化日志实例 |
BatchNum 和 SyncDuration 任一条件先满足即触发批量写入。
writer.ESMapping 定义 MySQL 表到 ES 索引的映射关系:
| 字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
TableName |
string |
是 | 数据库表名 |
Index |
string |
是 | ES 索引名 |
FilterColumns |
[]string |
否 | 需要同步的列(白名单),为空则同步全部 |
IgnoreColumns |
[]string |
否 | 忽略的列(黑名单) |
Handler |
func(data map[string]any) bool |
否 | 行级过滤器,data 为已格式化后的列名→列值映射,返回 false 则跳过该行 |
FilterColumns和IgnoreColumns同时设置时,仅FilterColumns生效。Handler 收到的data已经过列过滤和类型格式化处理。
MySQL binlog 事件
│
▼
┌──────────┐ ┌──────────────────┐
│ Canal │────▶│ writer.Writer │
│ (binlog) │ │ (数据写入目标) │
└────┬─────┘ └────────┬─────────┘
│ │
▼ ▼
┌──────────┐ ┌──────────────────┐
│ Channel │ │ position.Service │
│(容量2048)│ │ (同步进度管理) │
└────┬─────┘ └────────┬─────────┘
│ │
▼ ▼
┌──────────┐ ┌──────────────────┐
│ consume │ │PositionRepository│
│(批量聚合)│ │ (位点持久化) │
└──────────┘ └──────────────────┘
核心调度器(内部类型),负责:
- 通过 go-mysql 的 canal 组件伪装为 MySQL Slave,监听 binlog
- 首次启动时对未 dump 的表执行全量数据 dump
- 将 binlog 事件通过 channel 送入消费者 goroutine
- 按
BatchNum或SyncDuration条件批量调用 Writer 写入 - 每次写入成功后通过
position.Service持久化同步进度
type Writer interface {
Write(events []*event.Event) error
Tables() []string
}Write:批量写入事件到目标端Tables:返回该 Writer 关注的所有表名,供 syncer 过滤 binlog 事件
内置实现:writer.ESWriter(通过 NewESWriter 创建),使用 ES Bulk API 批量写入。INSERT 生成 index 操作,UPDATE 自动处理主键变更(先删旧文档再索引新文档),DELETE 生成 delete 操作。
管理 binlog 同步进度(当前 binlog 文件名、位点、已 dump 的表列表),确保重启后从上次中断的位置继续。每次批量写入成功后自动更新位点并持久化。
type PositionRepository interface {
Get(ctx context.Context, database string) (*Position, error)
Set(ctx context.Context, database string, pos Position) error
}定义同步进度的持久化层接口。内置 ES 实现(通过 NewESRepo 创建),将进度存储在 .sync_position 索引中。实现此接口即可替换为 Redis、etcd、数据库等存储后端。
封装单个 binlog 行事件:
type Event struct {
LogFileName string // binlog 文件名
RowsEvent *canal.RowsEvent // 原始 binlog 行事件(包含 Table、Action、Rows 等信息)
}event.FormatValue(column, value) 工具函数将 MySQL 列值转换为适合 JSON 序列化的格式,完整支持 DATETIME、TIMESTAMP、DATE、TIME、ENUM、SET、BIT、JSON、BLOB、BINARY/VARBINARY 等类型。
Handler 在写入前对行数据进行二次过滤,签名为 func(data map[string]any) bool:
writer.ESMapping{
TableName: "orders",
Index: "orders_index",
Handler: func(data map[string]any) bool {
// 只同步 status 为 "active" 的订单
status, ok := data["status"].(string)
return ok && status == "active"
},
}实现 position.PositionRepository 接口即可切换存储后端:
// 示例:使用 Redis 存储同步进度
type redisRepo struct {
client *redis.Client
}
func (r *redisRepo) Get(ctx context.Context, db string) (*position.Position, error) {
// 从 Redis 读取进度...
}
func (r *redisRepo) Set(ctx context.Context, db string, pos position.Position) error {
// 写入 Redis...
}
// 使用时
ps, _ := position.NewService(&redisRepo{...}, "mydb", slog.Default())
syncer, _ := mysqlsyncer.NewSyncer(dsn, esWriter, ps)实现 writer.Writer 接口可同步到任意目标端:
type kafkaWriter struct {
producer *kafka.Producer
tables []string
}
func (w *kafkaWriter) Write(events []*event.Event) error {
// 将事件序列化后发送到 Kafka...
}
func (w *kafkaWriter) Tables() []string {
return w.tables
}
// 使用时
kw := &kafkaWriter{tables: []string{"users", "orders"}}
ps, _ := position.NewServiceWithES(esClient, "mydb", slog.Default())
syncer, _ := mysqlsyncer.NewSyncer(dsn, kw, ps)本项目基于 MIT License 开源。