ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

AWS MySQL binlog数据同步实战:从方案选型到踩坑记录

AWS MySQL binlog数据同步实战:从方案选型到踩坑记录 先说结论这个项目的核心链路不复杂就是把 AWS 上 MySQL/RDS 的 binlog 打开用解析工具把 row 格式的变更事件读出来再经过消息队列喂给下游存储或业务系统。我最初接到这个需求的时候想法很简单“解析 binlog 嘛网上不是一堆方案吗”结果一上手才发现坑全藏在细节里尤其是 AWS 托管数据库和自建 MySQL 在参数、权限、位点管理上有不少差异不提前摸清楚后面同步跑起来就是各种断流、丢数据、追不上。这篇文章我会把我在 AWS 环境下做 MySQL binlog 数据同步的完整方案、选型理由、配置过程、遇到过的坑都梳理出来给准备做 CDC 实时同步的朋友一份可以直接抄作业的参考。1. 方案选型为什么最终选了解析 binlog1.1 先对比一下常用的数据同步方式当时业务需求是把线上订单表实时同步到 Elasticsearch提供搜索和聚合查询。第一版想得简单直接用定时任务跑增量查询每次取 update_time 大于上次轮询时间的数据同步到 ES。这个方案对数据量小、删除场景少的业务确实够用但问题也很明显删除操作天然无法感知只能用软删除标记一旦某个模块忘了做标记ES 里就多出一堆脏数据另外每次轮询都会打到源库的索引表大了以后更新频繁update_time 的索引会拖慢写库。后来又考虑了应用双写也就是下单接口在写 MySQL 的同时手动写一份到 ES。这个方案看似可控实际是给业务代码绑了一颗定时炸弹所有改库的地方都得跟着改一旦 ES 抖动主流程跟着失败运维同学半夜起来捞日志的场景想想就头疼。触发器方案就更不建议了云数据库上维护触发器本来就不方便还会给源库增加额外开销属于性价比最低的选择。最后聚焦到解析 binlog 做增量同步。这个方案并不新鲜本质就是模拟 MySQL 主从复制里的从库角色读取主库的 binary log 事件流然后把每一行变更转成结构化数据。它的优点一个是实时性高事件产生后毫秒级就能被消费另一个是对业务完全无侵入不需要改任何业务代码最关键的是 row 格式的 binlog 里记录了每一行的完整变更镜像插入、更新、删除都能拿到这是其他方案都做不到的。方案实时性对业务侵入捕获删除实现复杂度定时轮询 update_time分钟级低不支持低应用双写实时高支持中触发器中间表准实时中支持中解析 binlog实时低支持高1.2 binlog 的三种格式为什么必须用 rowMySQL 的 binlog 有三种格式statement、row、mixed。statement 格式记录的是原始 SQL比如UPDATE orders SET status1 WHERE id100看起来体积小但实际上在同步场景里很难用因为你不知道这条 SQL 影响了哪些行还得去源库重新执行才能拿到结果而再执行一次又可能产生新的数据变化根本没法保证一致性。mixed 格式是 MySQL 根据语句自动选择 statement 或 row但这种不确定性的东西用在同步链路上就是灾难你不知道拿到的 event 到底带不带行数据。所以做数据同步必须把 binlog_format 强制设为 row。row 格式下每个事件都包含了变更行的完整数据快照插入事件里有整行数据更新事件里有变更前和变更后的镜像删除事件里有被删掉的整行。解析端拿到这些数据不需要再回源库查直接就能转化、分发这也是为什么所有 CDC 工具都要求源库开启 row 格式。在 AWS 的 RDS for MySQL 上光设置 binlog_formatROW 还不够我强烈建议把 binlog_row_image 也显式设置成 FULL这样才能保证更新事件里携带完整的前后镜像。如果用的是默认值或者 MINIMAL更新事件里可能只有主键和发生变化的那几个字段这对下游需要全量数据的场景会很麻烦。结合我自己的经验这两个参数一个都不能省后面踩坑部分会细说。1.3 AWS 场景下为什么没选 DMS既然人在 AWS 上可能有人会问为什么不直接用 AWS Database Migration Service 的持续复制功能DMS 确实能把 RDS MySQL 的数据持续复制到 S3、Redshift 等目标但它的目标是“数据库迁移”和“简单同步”对数据转换能力有限而且目标端如果想接自建的 Elasticsearch、Redis 或者业务自定义格式DMS 做起来很别扭。自己解析 binlog 的好处是可以拿到最原始的变更事件流想怎么加工怎么加工转成 JSON 扔给 Kafka、SQS、或者直接调接口都行灵活性完全在自己手里。一句话总结我的选型逻辑如果目标是简单的库到库迁移选 DMS如果目标是构建一条属于自己的实时数据管道让下游能自由消费订单变更那就自己解析 binlog。我这次的需求明显是后者。2. AWS 上把 binlog 打开并准备同步环境2.1 RDS 参数组配置改完要重启我用的是 RDS for MySQL 8.0默认情况下 binlog_format 是 MIXED虽然也会产生 binlog但格式不是我们想要的。要改的话需要到 RDS 控制台找到实例关联的参数组修改下面几个关键参数binlog_format ROWbinlog_row_image FULLmax_binlog_size 128默认一般就是 128MB保持默认即可binlog_expire_logs_seconds 按需设置这是 MySQL 8.0 控制 binlog 保留时间的参数需要注意修改参数组之后需要重启 RDS 实例才能生效。AWS 控制台会提示“pending reboot”重启会造成几十秒到几分钟的连接中断如果业务不能接受得提前安排维护窗口。这里有一个容易忽略的地方如果这个参数组被多个实例共用修改后所有关联实例全部受影响建议单独为同步项目建一个专属参数组。提示RDS 上有个自动备份参数叫 binlogBackup默认开启它会把 binlog 记录到 S3 做备份这个不影响我们解析线上 binlog不需要动。2.2 binlog 保留时间一定要手动调长这是我在 AWS RDS 上吃过大亏的地方。自建 MySQL 可以通过 expire_logs_days 或 binlog_expire_logs_seconds 控制 binlog 清理时间但 RDS for MySQL 提供的是另一个存储过程接口CALL mysql.rds_set_configuration(binlog retention hours, 72);这条语句会把 binlog 在实例上保留 72 小时。为什么不建议用参数组里的 binlog_expire_logs_seconds因为 RDS 上该参数的默认逻辑会被rds_set_configuration覆盖官方也推荐用这个存储过程来管理所以别自己去改参数组里的过期时间容易产生配置冲突。我当时上线初期只设置了 24 小时结果有一次消费程序夜里出故障第二天早上才发现想再从 binlog 拉数据已经晚了凌晨那段 binlog 被清掉了最后只能重新做全量增量初始化非常痛苦。现在我给所有同步链路的 binlog 保留时间都设为 72 小时起步宁可多占点存储也要给自己留出故障恢复的窗口。查询当前配置SELECT * FROM mysql.rds_configuration;2.3 创建最小权限账号别用 root 去解析解析 binlog 需要一个账号但不需要所有权限。按最小权限原则创建专门给同步链路用的账号CREATE USER binlog_sync% IDENTIFIED BY your_strong_password; GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO binlog_sync%;REPLICATION SLAVE 权限让账号可以模拟从库拉取 binlog 事件流REPLICATION CLIENT 权限允许账号执行 SHOW MASTER STATUS、SHOW BINARY LOGS 等命令用于查询当前位点。如果解析工具需要读表结构做字段映射那就再给对应业务库的 SELECT 权限GRANT SELECT ON your_db.* TO binlog_sync%;安全组方面如果解析程序部署在 EC2 上建议只放行该 EC2 私网 IP 访问 RDS 的 3306 端口不要对全公网开放。我见过有人贪图方便把 RDS 安全组设成 0.0.0.0/0结果被扫描爆破这个是真不该省。2.4 本地测试环境用 Docker 快速验证生产环境没准备好之前建议本地用 Docker 起一个 MySQL 实例模拟 binlog 输出解析脚本可以先在本地跑通再切到 AWS。启动命令很简单docker run -d --name mysql-binlog-test \ -e MYSQL_ROOT_PASSWORDroot \ -p 3306:3306 \ mysql:8.0 \ --server-id1 \ --binlog-formatROW \ --binlog-row-imageFULL \ --gtid-modeON \ --enforce-gtid-consistencyON这样本地就有了一个外网可访问的 MySQL插入一条数据然后写个解析脚本消费一下就能完整模拟“源库产生变更 → 解析端拉取 → 输出 JSON”的流程比直接在 RDS 上调试效率高很多。3. 解析工具选型与核心配置3.1 主流 binlog 解析工具怎么选工具选型看得不是谁功能多而是和自己的技术栈、团队运维能力匹配。我列一下现在市面上用得多、口碑也比较稳的几款。Canal 是阿里开源的老牌产品Java 技术栈能够伪装成 MySQL 从库拉取 binlog支持 TCP、Kafka、RocketMQ 几种投递模式。它最大的优势是社区成熟文档多生产环境验证充分缺点是部署有点重包含 canal server、canal adapter、canal admin 好几个组件小团队上来会被绕晕。Maxwell 是轻量级的守护进程基于 Java 写的配置简单一条启动命令就能跑默认把 binlog 事件解析成 JSON可以直接输出到 stdout 或 Kafka。它很适合快速落地但复杂数据处理和灵活的路由能力偏弱。Debezium 是红帽开源的项目架构在 Kafka Connect 之上和 Kafka 生态绑定很深。如果你的基础设施里已经有 Kafka并且团队熟悉 Kafka Connect它是非常强大的选择但如果没有 Kafka为了一个同步任务去部署一套 Kafka Connect投入产出比就不太划算。python-mysql-replication 不是一个独立服务而是一个 Python 库非常轻量适合写脚本嵌入式集成到自己的同步程序里。我这次的项目量级在每秒几百条以内目标端是 Elasticsearch 和 Redis用这个库完全够用。工具技术栈部署复杂度适合场景CanalJava中高大规模同步、需要投递到 Kafka/RocketMQMaxwellJava低轻量 JSON 输出快速落地DebeziumJava/Kafka Connect高已有 Kafka 生态python-mysql-replicationPython低中小业务、嵌入式开发3.2 Canal 快速部署与配置示例如果你的同步链路非常长比如几十张表、下游还有多个系统我建议直接用 Canal省去自己造轮子。这里给一个 canal 1.1.7 的部署参考核心两步下载解压 canal deployer配置 canal.properties 和 instance.properties。canal.properties 中比较关键的几个配置canal.serverMode kafka canal.mq.servers your-kafka-broker:9092 canal.mq.retries 0 canal.mq.batchSize 16384instance.properties 指向源库canal.instance.master.address your-rds-endpoint.cluster-xxxx.rds.amazonaws.com:3306 canal.instance.dbUsername binlog_sync canal.instance.dbPassword your_strong_password canal.instance.connectionCharset UTF-8 canal.instance.filter.regex your_db\\..* canal.mq.topic binlog_topiccanal 会自动维护位点支持 failover比我们自己管理位点确实省心。但它默认会缓存最近一段时间的 binlog 位点如果 canal 停机时间超过保留窗口重启时同样会碰到“位点不存在”的问题这个逃不掉binlog 保留时间依然得设置得足够长。3.3 用 Python 写一个最轻量的兼容方案我不太喜欢把工具选得太重而且这次项目需要跟公司内部的告警系统打通用 Python 库反而更容易集成。python-mysql-replication 这个库在 PyPI 上叫 mysql-replication安装pip install mysql-replication核心代码特别简单先定义连接信息创建 BinLogStreamReader然后遍历事件判断事件类型把行数据转成 JSONfrom pymysqlreplication import BinLogStreamReader from pymysqlreplication.row_event import WriteRowsEvent, UpdateRowsEvent, DeleteRowsEvent import json mysql_settings { host: your-rds-endpoint.cluster-xxxx.rds.amazonaws.com, port: 3306, user: binlog_sync, passwd: your_strong_password, charset: utf8mb4, } stream BinLogStreamReader( connection_settingsmysql_settings, server_id123456, blockingTrue, only_events[WriteRowsEvent, UpdateRowsEvent, DeleteRowsEvent], resume_streamTrue, log_filemysql-bin-changelog.000001, log_pos4, ) for event in stream: for row in event[rows]: record { db: event.schema, table: event.table, event_type: event[event_type].decode() if isinstance(event[event_type], bytes) else event[event_type], data: row, log_file: event.packet.log_file, log_pos: event.packet.log_pos, } # 这里把 record 发到 SQS / Kafka / 直接写 ES按需实现 print(json.dumps(record, ensure_asciiFalse, defaultstr)) stream.close()注意 server_id 不能和源库上已有的从库、其他解析实例重复否则 MySQL 会认为这是一个冲突的从库连接把其中一个连接踢掉。我曾经在同一套环境上跑两个解析脚本server_id 都是 123456结果一个频繁断连日志里全是访问权限报错研究半天才发现是这个原因。另外这个库还支持只监听特定表、过滤数据库等比如增加 only_schemas、only_tables 参数。我们一般用 only_tables 把同步的表白名单列出来避免收到无关变更。4. 从 binlog 到目标存储同步链路怎么搭4.1 整体架构解析端和消费端要分离很多人在第一次做 binlog 同步时会犯一个错误在解析事件的同时直接同步去写目标库。这样解析逻辑和业务逻辑高度耦合一旦目标端抖动解析进程会被阻塞binlog 消费速度一掉源库磁盘占用就会持续上涨。我建议的架构是加一层消息队列解析端只负责读取 binlog 事件把结构化 JSON 发到消息队列然后由独立的消费端从队列里拉取数据写入目标存储。这让我在同步链路的中间位置有了缓冲能力我们可以控制消费速度、做失败重试、甚至把消息再投递给其他系统。在 AWS 上如果团队对 Kafka 运维不熟可以先用 SQS 顶一顶数据量上来再切 Kafka。完整的链路是这样的RDS MySQL 产生 binlog → 解析程序伪装成从库拉取 row 事件 → 序列化为 JSON → 发送到 Kafka/SQS → 消费程序读取消息 → 根据表名路由到目标 ES/Redis → 写入成功后提交位点。4.2 位点管理断了怎么续这个必须想清楚binlog 的事件流是顺序且连续的每个事件都对应一个 binlog 文件名加一个 position 偏移量这就是位点。解析端必须记住自己消费到了哪个位点否则程序重启后要么丢数据要么重复消费。python-mysql-replication 里resume_streamTrue 配合 log_file、log_pos可以指定从某个位点开始读取。一个稳妥的做法是把每次解析到的位点记录到本地文件、数据库或状态表中但不要每一条都写一次那样写频繁太高。我通常的做法是批量消费处理完一批后用最后一条事件的位点更新状态存储。伪代码逻辑是这样的启动时读取本地保存的位点log_file log_pos 如果本地有位点就从该位点启动 如果没有位点就从当前主库位点SHOW MASTER STATUS启动 每处理完一批事件更新本地位点 程序崩溃重启后从上次成功提交的位点继续这里要特别注意如果你每处理完 1000 条或每 1 秒更新一次位点那么在“刚消费完一批、还没来得及更新位点”这个窗口里崩溃重启后会有少量重复消费。下游写入必须做幂等这是必须配套的下一节细说。4.3 全量增量衔接防止丢数据增量同步之前通常要先做一次全量数据初始化把存量数据导入目标端然后再启动 binlog 增量消费。如果顺序搞反了就会出现存量导入覆盖增量变更、或者漏掉初始化期间产生的新数据的问题。推荐的流程是在源库执行SHOW MASTER STATUS;记录当前 binlog 文件名和 position。用 mysqldump 做全量导出mysqldump -h your-rds-endpoint -u binlog_sync -p \ --single-transaction --set-gtid-purgedOFF \ your_db your_db.sql--single-transaction 是利用 InnoDB 的一致性快照实现非锁表导出RDS 上推荐这样用不影响线上写入。全量导出的文件导入目标端。从第一步记录的位点开始启动增量解析程序。注意第 2 步的导出时间和第 1 步的记录时间之间如果有写入这些写入事件对应的 binlog 位置一定在记录的位点之后所以增量从那个位点开始不会丢数据。不要反过来先导入全量再记录位点那样中间的数据就丢了。4.4 下游写入的幂等与乱序处理binlog 事件本身是顺序的但加入消息队列和消费并发后顺序就不一定能保证了。比如同一行在主库先 UPDATE 后 DELETE两个事件进了不同分区或不同消费线程目标端可能先看到 DELETE再看到 UPDATE最后把一条不存在的记录又写回去了。我的经验是两条第一尽量按表名或主键哈希决定消息分区保证同一主键的消息进同一个消费线程第二消费端写入必须幂等用目标存储的天然更新语义来覆盖。比如写 Elasticsearch 时用业务主键作为文档 _id写 Redis 时用 SET 覆盖写关系型库时用 INSERT ... ON DUPLICATE KEY UPDATE这样即使重复消费、乱序消费最终也是最后一次事件生效。如果两条更新事件本身就存在先后顺序下游存储不支持版本控制建议在消息里带上事件产生时间或 binlog 位点消费端根据位点大小决定是否覆盖避免旧事件把新事件覆盖掉。5. 整个方案跑起来后我踩过的那些坑5.1 经典问题和排查思路先列一张速查表这些都是我在实际项目中遇到过的如果你也做 binlog 同步大概率会碰到。现象可能原因解决办法程序启动报 impossible positionbinlog 已被清理保存的位点太旧延长 binlog 保留时间重新做全量增量解析进程频繁断连server_id 与源库其他从库冲突修改 server_id 为全局唯一值同步延迟逐渐拉大消费能力不足binlog 产生速度大于消费速度批量消费、增加并发、优化目标端写入中文乱码连接字符集没配对连接参数加 charsetutf8mb4更新事件里没有完整镜像binlog_row_image 不是 FULL参数组修改后重启实例目标端数据错乱乱序消费或重复消费未做幂等按主键分区、upsert 写入RDS 磁盘空间一直涨binlog 保留时间过长且消费滞后提高消费速度或缩短保留小时数5.2 字符集和时区这两个小坑MySQL 的事件数据里字符串字段的字节码是跟着源库 charset 走的如果解析端连接时用的 charset 和源库不一致就会导致中文乱码。在 Python 的 mysql-replication 库中务必在 mysql_settings 里配置charset: utf8mb4。如果源库有些表还是老旧的 latin1就要在连接配置里用源表的真实字符集否则解析出来的数据会变成问号。时区问题更隐蔽。MySQL 的 timestamp 类型在 binlog 里存储的是 UTC 时间datetime 类型存的是字面时间不带时区概念。如果业务存储在 datetime 字段里解析端拿到的字符串就是原生值这个没问题但如果是 timestamp不同连接时区下解析出的时间可能不一样。我的做法是统一规范上游建议业务尽量用 timestamp解析端和消费端统一使用 UTC 时间处理展示层再转本地时区这样最不容易出歧义。5.3 性能优化和监控告警性能方面解析脚本处理单个事件通常很快瓶颈基本都在目标端写入。批量写入的收益非常明显比如写 Elasticsearch 时用 bulk 接口一批几百条提交一次比逐条写入快不止一个量级。但批量也不能太大否则单次请求失败重放的成本太高一般控制在 500 到 1000 条一提交比较合适。监控是我强烈建议大家一上来就做好的东西不要等出了问题再去补。在 AWS 上RDS 控制台本身有 BinLogDiskUsage 指标可以看 binlog 占用的存储空间这个指标一旦持续攀升说明消费速度跟不上了。我自己还会写一个脚本定时执行SHOW MASTER STATUS把当前位点和解析程序上报的已消费位点做差值换算成滞后秒数发到 CloudWatch。滞后时间超过阈值就告警这个比看磁盘指标更直接。消费端建议开启失败重试和死信队列。如果一条消息反复处理失败不要无限重试应该进死信队列保留下来人工排查后再回放。比如 Elasticserch 索引 mapping 冲突这一类问题重试一百遍也没用必须人工介入改 mapping 才能解决。5.4 上线前的自我检查清单我把这些经验整理成一个 checklist每次新接同步任务我都会按这个过一遍binlog_format 是否为 ROWbinlog_row_image 是否为 FULLbinlog 保留时间是否足够建议首次上线阶段设置到 72 小时以上解析账号是否最小权限密码是否已托管到密钥管理服务安全组是否只放行必要的来源 IP解析端和 RDS 是否同一 VPCserver_id 是否全局唯一有没有和现有从库冲突位点是否有持久化机制程序重启后能否正确续跑全量初始化流程是否验证过是否清楚记录到增量启动的衔接位点下游写入是否幂等目标端主键/文档 id 是否统一CloudWatch 指标与告警是否配置滞后多少秒会触发通知这些事项看着琐碎但每一条背后都有真实的故障案例支撑。我在上线初期吃过 binlog 保留时间太短的亏、也吃过安全组暴露的亏后来把清单固化下来新项目照着走踩坑概率大幅下降。我在实际项目中最后用的是 python-mysql-replication 配合 SQS 消费者这套方案运行了小半年数据量不大但很稳定。如果你问我有什么建议我最大的体会是不要一上来就把链路搭得太重先摸清自己的数据量、消费速度、下游能力再决定用 Canal 还是轻量脚本。另外再分享一个小技巧刚开始跑的时候把 binlog 保留时间设置到 72 小时然后每天看一眼消费延迟曲线等连续一周没有异常了再慢慢把保留时间缩短到 24 或 48 小时。这样既不会因为配置太短丢失位点也不会因为配置太长白白占用存储两全其美。
返回列表