异构数据源同步必须按源类型分治:MySQL用binlog、PostgreSQL用logical replication、ES用_version+refresh、MongoDB用change stream;位点不可互换,需独立存储(如SQLite),避免强耦合导致同步卡死。

异构数据源同步不能靠统一 SQL 或通用 ORM 解决,必须按源类型分治:MySQL 用 binlog、PostgreSQL 用 logical replication、ES 用 _version + refresh 控制、MongoDB 用 change stream —— 混用一套逻辑必然丢数据或重复。
MySQL 和 PostgreSQL 增量位点不可互换
MySQL 的 binlog position(文件名+偏移)和 PostgreSQL 的 LSN(64 位整数)本质不同,不能简单存成同一个字段再“统一解析”。go-mysql-org/go-mysql 启动时传入的是 mysql.Position{File: "mysql-bin.000001", Pos: 4},而 pglogrepl.StartReplication 要的是 pglogrepl.LSN(0x12345678)。两者无法对齐时间语义,更不能跨库比较大小。
- MySQL 位点必须配合
binlog_format = ROW,否则RowsEvent拿不到具体字段变更 - PostgreSQL 必须提前建 slot:
SELECT * FROM pg_create_logical_replication_slot('my_slot', 'wal2json');,且wal_level = logical要写进postgresql.conf - 别把 MySQL 的
server_id和 PG 的 replication user 权限混管——前者是连接标识,后者需REPLICATION角色
ES 同步必须显式控制 refresh 和 version
直接调 client.Bulk() 写入 ES,文档可能几秒后才可查,且副本延迟会导致下游读取为空;若没校验 _version,并发更新会覆盖而非报错。ES 7+ 默认禁用外部 version 检查,不手动开启就等于放弃幂等性。
- 建索引时加配置:
"version_type": "external",否则version字段被忽略 - 写入时带 version:
client.Index().Index("my-index").Id("123").Version(100).VersionType("external") - 调试可用
Refresh("true"),但生产环境 bulk 大小应 ≤ 1000 条,避免触发es_rejected_execution_exception - 别依赖 ES 自身的
_seq_no做同步位点——它不跨分片单调,也不保证全局顺序
MongoDB change stream 依赖 oplog 和 replica set
单节点 MongoDB 不支持 change stream,必须部署为 replica set(哪怕只一个节点),且 oplog 大小要足够容纳同步延迟窗口。Golang 官方驱动的 Collection.Watch() 返回的是 changeStream,但事件结构松散,FullDocument 默认不返回,UpdateDescription 字段需手动解析。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
立即学习“go语言免费学习笔记(深入)”;
- 启动前确认:
rs.initiate()已执行,且oplogSizeMB≥ 1024(尤其写入密集时) - Watch 时加
FullDocument: options.ChangeStreamFullDocumentUpdateLookup,否则update事件只有updatedFields,无旧值 - 不要用
ResumeToken做持久化位点——它不跨进程稳定,应转成clusterTime或documentKey._id落盘 - change stream 断连后默认重试,但超时时间由
MaxAwaitTime控制,建议设为 30s 防止长阻塞
同步位点存储必须与目标库解耦
把 MySQL 的 binlog pos 存进 PostgreSQL、把 PG 的 LSN 存进 ES、把 MongoDB 的 resume token 存进 Redis —— 这些都属于“强耦合”,一旦目标库不可用,位点就无法更新,整个同步卡死。真正可靠的方案是用独立、轻量、本地可落盘的存储,比如 SQLite 文件或本地 JSON 文件。
- 位点表/文件必须带写锁或原子 rename,防止多进程同时写坏
- 每次成功提交一批变更后,再原子更新位点;失败则原位点保留,下次重拉
- 别在事务里写位点(如 BEGIN → 写 PG → 写位点 → COMMIT)——PG 提交成功但位点写失败,就会重复消费
- SQLite 是最省心的选择:
CREATE TABLE sync_cursor (source TEXT PRIMARY KEY, cursor TEXT NOT NULL);,单文件、零依赖、ACID
异构同步最难的不是连上每个库,而是让它们的“进度”在逻辑上对齐:MySQL 的 position、PG 的 LSN、Mongo 的 clusterTime、ES 的 timestamp,没有一个能天然互转。你得接受它们就是不同世界的刻度,只能靠业务语义(比如订单创建时间)做软对齐,或者干脆放弃全局一致,按源独立推进。

















