直接用golang调ES API同步易丢数据,因ES写入非强一致、副本延迟且程序常未等refresh或校验_version,重启时又未持久化消费位点;需用可靠队列+幂等写入+位点落盘保障at-least-once。

为什么直接用 golang 调 ES API 做同步容易丢数据
因为 ES 的写入不是强一致,index 或 bulk 成功只代表主分片写入成功,副本可能延迟;而 golang 程序如果没等 refresh 或没校验 _version,后续读取就可能查不到刚写入的文档。更常见的是,程序重启时没保存消费位点(比如 MySQL binlog offset 或 Kafka offset),导致重复或漏同步。
用 elastic/v7 客户端做 bulk 同步的关键配置
官方客户端 elastic/v7 默认不启用 refresh,也不自动重试失败的 bulk item,必须手动控制:
-
client.Bulk().Index("my-index").Refresh("true"):强制刷新,让文档立即可查(仅用于调试或低吞吐场景,生产慎用) - 务必检查
res.Errors == true,并遍历res.Items逐个判断每个 action 的Status和Error.Reason - 设置
client.SetSniff(false)避免节点发现干扰长连接,配合SetHealthcheck(false)减少心跳开销 - bulk 大小建议控制在
1000条以内、单条≤ 1MB,否则易触发 ES 的es_rejected_execution_exception
如何保证 MySQL → ES 的 at-least-once 同步
核心是把「数据变更」和「ES 写入成功」绑定到同一个事务边界——这在跨系统时做不到,所以退而求其次:用可靠队列 + 幂等写入 + 位点持久化:
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 监听 MySQL binlog(用
go-mysql-elasticsearch或自研 canal client),把变更解析为id+op+json发到 Kafka/RocketMQ - golang 消费者用
consumer.CommitOffsets()手动提交 offset,且只在 ESbulk全部成功后才提交 - ES 写入前加
if_seq_no和if_primary_term防覆盖;或用upsert+ 脚本判断字段是否已存在 - 位点存储不能只放内存,必须落盘(如写入本地
offset.json或 PostgreSQL 表),且每次 commit 后fsync
遇到 version_conflict_engine_exception 怎么办
这不是 bug,是并发更新冲突的正常反馈。直接忽略或重试会掩盖问题,正确做法是结合业务语义处理:
立即学习“go语言免费学习笔记(深入)”;
- 如果是「最终一致性」场景(如商品库存摘要),可忽略该错误,下次变更自然覆盖
- 如果是「强顺序依赖」(如用户积分变更),需在写入前先
GET当前_seq_no,再带条件写入;失败则重拉最新版本再计算 - 避免用
update接口做高频计数器累加,改用script+params原子操作,例如:"source": "ctx._source.count += params.delta" - 注意:ES 7+ 默认禁用
_version检查,如需启用,建索引时设"version_type": "external"
同步链路里最常被跳过的其实是错误回溯能力——没有记录每条失败文档的原始 payload 和 ES 返回体,出问题只能盲猜。留好 failed_items.log 比调参重要得多。

















