go-mysql-elasticsearch 不适合自定义同步逻辑,因其将 binlog 解析、es 写入、位点管理强耦合,修改字段映射、过滤、索引切分或对接 kafka 需 fork 硬改,升级困难;应选用可插拔的 go-mysql 库,它提供稳定 row-based event 拉取与原始 event 结构体暴露。

为什么直接用 go-mysql-elasticsearch 不适合自定义同步逻辑
它把 MySQL binlog 解析、ES 写入、位点管理全打包了,一旦你要改字段映射、加过滤条件、切分索引或对接 Kafka,就得 fork 代码硬改,后续升级困难。真正需要的是可插拔的 binlog 解析层 —— mysql-binlog-event 或更主流的 github.com/go-mysql-org/go-mysql。
后者自带 Syncer 和 BinlogStreamer,能稳定拉取 row-based event,并暴露原始 Event 结构体,是你做定制同步的起点。
-
go-mysql默认只支持 ROW 格式 binlog,确保 MySQL 开启binlog_format = ROW,否则会报Unknown binlog format - 不要用
mysqldump或SELECT ... FOR UPDATE混合事务,这类操作不会生成 row event,导致数据丢失 - 位点(
Position)必须持久化到外部存储(如 Redis 或本地文件),不能只存在内存里,否则进程重启就从头重放
如何解析 Delete/Update/Insert event 并提取变更字段
row event 的核心是 RowEvent,但它不直接暴露字段名,需结合 TableMapEvent 才能还原列名和类型。常见错误是直接遍历 Rows 字节切片,结果解出乱码或 panic。
正确做法是:先监听 TableMapEvent 缓存表结构(schema.table → columnNames, columnTypes),再在后续 WriteRowsEvent/DeleteRowsEvent/UpdateRowsEvent 中复用该映射。
-
UpdateRowsEvent的Rows字段是[][]interface{},每对元素为[old_values, new_values];DeleteRowsEvent只有old_values;WriteRowsEvent只有new_values - 注意
NULL值:go-mysql 用nil表示,不是sql.NullString,判空直接用v == nil - 时间类型(
DATE/DATETIME)默认解析成[]byte,需手动转time.Time,例如:time.Parse("2006-01-02 15:04:05", string(b))
如何安全地推进 binlog 位点(GTID vs File/Pos)
用 GTID 更可靠,但要求 MySQL 5.7+ 且 gtid_mode = ON。如果环境不支持,只能退回到 file + position 模式,此时必须保证「解析 → 处理 → 提交位点」是原子的,否则会丢数据或重复。
推荐方案:把位点更新和业务写入放在同一个事务里(比如用 PostgreSQL 做同步元数据表),或至少用两阶段提交模拟 —— 先写下游成功,再更新位点。别用「先更新位点再写下游」,那是灾难配置。
- 使用
Syncer.StartFrom()时,传入的Position必须是已存在的 binlog 文件,否则报Could not find first log file name in binary log index file - GTID 模式下,
Syncer的StartFromGTID()需传入mysql.GtidSet实例,不能传字符串;可用mysql.ParseGtidSet("de278ad0-xxxx-11e4-b3d5-0025904b832c:1-100") - 每次成功处理一个 event 后,调用
syncer.SetNextPosition(pos)更新内存位点,但一定要配合外部持久化,否则崩溃即丢失
为什么解析后直接写 ES/Kafka 容易背压崩溃
MySQL 产生 binlog 的速度远高于下游写入吞吐,尤其批量 delete 或 alter table 导致单秒数千 event。若用同步阻塞方式逐条写 ES,Syncer 内部 channel 会迅速填满,最终触发 syncer: event channel is full panic。
必须引入缓冲与降级机制:用带界线的 chan *RowEvent(如 1000 容量),配合 goroutine 消费;当 channel 满时,暂停 Syncer 的 event 拉取(调用 syncer.Pause()),等消费进度追上再 Resume()。
- 别用
log.Fatal或 panic 处理写入失败,应重试 + 降级到本地磁盘暂存(如按时间分片写/tmp/binlog_backlog_20240415.log) - Kafka 场景下,event 的 key 应设为
schema_table_pk(如"user_123"),保证同一行变更顺序写入同一 partition - ES bulk 写入建议 100–500 条/批,太小则 HTTP 开销大,太大则单次失败影响面广
binlog 解析真正的难点不在读,而在「状态一致性」—— 位点、下游写入、本地缓存三者任何时候都不能错位。哪怕加一行日志,也要想清楚它发生在位点更新前还是后。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!











