如何用Python实现MySQL到Elasticsearch的新增数据同步最优方案
MySQL增量数据同步至Elasticsearch最优实现方案
原有双状态列方案的修复(最小改造成本)
原有方案重复读的核心原因是查询和更新操作非原子,多实例同时查询到未处理行后,都会触发更新逻辑。最小改动修复方式如下:
- 替换「先查后改」逻辑为原子抢占更新,先执行带行锁的UPDATE语句抢占待处理行,再查询已抢占的行做处理,确保同一行只会被一个实例抢到
- 参考SQL写法:
首先执行抢占更新:UPDATE 你的表名 SET process_start = 1 WHERE process_start = 0 AND process_done = 0 ORDER BY id LIMIT 1;
再查询已抢占的行:SELECT * FROM 你的表名 WHERE process_start = 1 AND process_done = 0 ORDER BY id LIMIT 1; - 新增超时重置逻辑:加
process_update_time字段,定时将process_start = 1且更新时间超过阈值(根据实际处理时长调整,建议设为正常处理时长的3倍)的行重置为未处理状态,避免实例中途挂掉导致行永久卡住 - 写入ES时用MySQL表的主键作为ES文档的
_id,就算出现重复处理也只会覆盖同一条文档,不会产生冗余数据
无侵入更优方案(无需修改原表结构)
方案1:Binlog监听同步(生产环境首选)
这是当前工业界最通用的MySQL增量同步方案,完全不侵入原有业务表:
- Python侧可以用
python-mysql-replication库监听MySQL的Binlog日志,直接获取实时的新增、修改、删除事件,解析字段后直接写入ES即可 - 多实例部署时可以将Binlog消费位点存在Redis等共享存储中,用分布式锁保证同一时间只有一个消费实例在运行,也可以将Binlog事件转发到Kafka等消息队列,实现多实例分布式消费
- 优势:延迟低(毫秒级感知数据变更)、对业务零侵入、没有额外的业务表写压力
- 前置要求:MySQL开启Binlog且格式设置为ROW模式,监听账号授予REPLICATION相关权限
方案2:基于自增主键/更新时间轮询(小数据量场景首选)
如果你的表自带自增主键或者update_time字段,不需要加任何额外字段即可实现:
- 自增主键方案:记录上一次拉取到的最大主键ID,每次轮询执行
SELECT * FROM 你的表名 WHERE id > 上一次最大ID LIMIT 1000;,写入ES完成后更新记录的最大ID即可,只适合只有新增、没有修改的业务场景 - 更新时间方案:适合有数据更新的场景,记录上一次拉取的时间戳,每次轮询执行
SELECT * FROM 你的表名 WHERE update_time > 上一次拉取时间戳 LIMIT 1000;,写入完成后更新时间戳 - 多实例部署时将拉取位点存在共享存储,加分布式锁保证同一时间只有一个实例轮询即可
内容的提问来源于stack exchange,提问作者Hemant Sah
相关产品推荐
相关产品推荐

