MySQL的ON DUPLICATE KEY UPDATE在大批量场景下会崩,因其单语句多行插入需逐行做唯一索引查找、加行锁、刷事务日志,10万行易致锁粒度失控、死锁或超时,且不支持WHERE过滤和删除语义,IO浪费严重,配合无超时控制的Go客户端更易阻塞goroutine。

直接说结论:大批量数据同步更新不能靠“全量覆盖”或“逐条 Upsert”,必须拆解为「位点驱动的增量流 + 幂等写入 + 分片缓冲」三层逻辑,否则 CPU、内存、数据库连接、网络带宽四项资源会同时打满。
为什么 MySQL 的 ON DUPLICATE KEY UPDATE 在大批量场景下会崩
它本质是单语句多行插入+冲突检测,但每条记录都要走一次唯一索引查找+行锁+事务日志刷盘。10 万行批量进来时:
- MySQL 内部会把这批语句拆成多个 mini-transaction,锁粒度失控,极易触发死锁或锁等待超时
-
INSERT ... ON DUPLICATE KEY UPDATE不支持WHERE条件过滤,无法跳过已同步完成的旧数据,白白消耗 IO - 如果源端有删除操作,该语法完全无法表达,只能靠额外 DELETE 语句补救,进一步放大事务体积
- Go 客户端用
db.Exec提交时,若没设context.WithTimeout,一个卡住的语句会让整个 worker goroutine 挂死
用 pglogrepl 或 go-mysql 拉 binlog/LSN 前必须确认的三件事
不检查这三项,同步启动即失败,或运行数小时后静默断连:
- MySQL:确认
binlog_format = ROW且server_id全局唯一;用SHOW VARIABLES LIKE 'binlog_format'和SHOW VARIABLES LIKE 'server_id'实时查,别信配置文件注释 - PostgreSQL:必须提前建好 replication slot,命令是
SELECT * FROM pg_create_logical_replication_slot('my_slot', 'wal2json');wal_level必须为logical,不是replica - 权限:MySQL 用户需有
REPLICATION SLAVE权限;PG 用户需有REPLICATION角色,且连接字符串里要显式加replication=database
bytes.Buffer 预分配容量不是“可选优化”,而是大批量同步的内存安全线
同步过程中常要把几百 KB 的 JSON event 批量序列化再发 Kafka 或写 ES,用默认 bytes.NewBuffer(nil) 就等于主动触发 OOM:
立即学习“go语言免费学习笔记(深入)”;
- 10 MiB 数据在未预分配时,
bytes.Buffer.grow可能执行 15+ 次扩容,每次复制已有数据,峰值内存占用翻倍 - 正确做法是根据单条 event 平均大小 × 批次数量估算,例如平均 2KB × 1000 条 = 2MiB,则初始化为
bytes.NewBuffer(make([]byte, 0, 2 - 更稳的方式是用
sync.Pool复用 buffer 实例,避免高频 GC;但注意 pool 中对象不能含跨 goroutine 引用(比如闭包捕获了 channel) - 千万别在 transformer 里用
fmt.Sprintf拼接大 JSON——它底层调bytes.Buffer且不预分配,比手写json.Marshal慢 3 倍以上
Sink 层写 Kafka/ES 时最容易被忽略的两个硬约束
它们不报错,但会导致数据丢失或乱序,且问题延迟暴露:
- Kafka:
sarama.SyncProducer默认 partitioner 是RandomPartitioner,相同 key 的 event 可能散落不同 partition,破坏变更顺序;必须显式设Partitioner: sarama.NewHashPartitioner,并确保每条 event 的msg.Key非空(比如用 user_id) - Elasticsearch:批量写入时若带
?refresh=true,每条 document 都会触发 segment 刷新,1000 条请求 = 1000 次 refresh,吞吐暴跌;应统一用?refresh=wait_for,让 ES 自行合并刷新时机 - checkpoint 存储介质必须和 sink 同可用区:Sink 写 us-east-1 的 Kafka,checkpoint 却存 us-west-2 的 S3,网络延迟高 + 跨区费用高 + 故障域不一致,一出问题就丢位点
真正卡住项目的从来不是“怎么读数据”,而是“怎么确认数据已稳稳落库且不可逆”。位点提交、buffer 生命周期、partition key 一致性、refresh 策略——这些细节没对齐,同步服务跑得越久,数据偏差越大。


















