ARTICLE DETAIL

资讯详情

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

MySQL到BigQuery数据同步:周期同步为何丢数据?CDC与binlog解析

MySQL到BigQuery数据同步:周期同步为何丢数据?CDC与binlog解析 如果你维护过一条 MySQL 到 BigQuery 的数据同步链路大概率有过这种经历每天凌晨跑一次定时任务把订单数据搬到数据仓库早上看报表一切正常。直到某一天业务方跑过来问为什么源库订单已经显示“已取消”BigQuery 报表里还是“待支付”。你查了同步脚本运行日志没有报错目标表也确实更新了但数据就是不对。如果这种问题反复出现你会慢慢意识到问题可能不在定时脚本本身而在“周期性同步”这个模式。周期性同步是大多数团队接触 MySQL 数据同步时的默认方案但这不代表它永远正确。当业务开始要求追溯变化、实时看板、审计历史时周期同步的短板会暴露得很明显。这也是 MySQL CDCChange Data Capture以及它背后的 binlog 越来越受关注的原因。这篇文章想讲清楚一件事周期同步到底漏掉了什么binlog 为什么能补上以及在真正落地到 BigQuery 之前你需要理解哪些过程和边界。1. 周期性同步真正漏掉的是什么1.1 快照只记录“当时是什么”不记录“发生了什么”周期同步不管你是做全量快照还是用updated_at 上次时间做增量本质上都把“数据同步”理解成了“把源表某个时间点的状态复制到目标表”。这样做没有错但你得到的只是一张最终状态的静态表。举个例子一条订单记录在一天之内经历了这样几次变化09:00 创建订单状态是“待支付”10:30 用户点击支付状态变成“已支付”11:00 风控触发异常状态变成“审核中”11:20 审核通过状态回到“已支付”如果每天凌晨同步一次BigQuery 里只能看到最后的“已支付”。中间那几次状态变化全部被覆盖。你可能会说正常日报只需要最终状态没问题。但一旦业务方想分析“从提交订单到完成支付平均要多久”或者“订单在审核状态停留了多长时间”这张周期同步表就回答不了。这不是因为你少写了什么判断逻辑而是周期性同步的模式天生不保留过程。过程被覆盖信息就丢失了。而 binlog 记录的是每一次变更事件它天然保留了这些过程。1.2 更新时间戳并没有想象中可靠很多团队不做全量快照而是用“增量同步”核心逻辑是SELECT * FROM orders WHERE updated_at 上次同步时间这个方案看起来很成熟但实际落地时坑很多。最常见的是应用层更新数据时没有更新updated_at字段。比如一个后台手动改单操作SQL 里只写了SET status 已取消没有带上updated_at NOW()。这时即使记录变了增量同步也完全感知不到。更隐蔽的情况是批量更新。比如运营工具把一批订单统一改成“已过期”可能直接执行UPDATE orders SET status expired WHERE expire_time NOW();如果这条 SQL 没有同步更新updated_at那么所有被更新记录对于增量同步来说都是“隐形变更”。还有一种常见情况两条记录在同一秒内被更新而时间戳精度只到秒。如果同步脚本按“大于上次同步时间”来抓可能漏掉边界。虽然可以通过偏移量或者毫秒精度来缓解但本质还是脆弱的。周期性增量同步过度依赖业务表里的某一列时间字段而这个字段既不能保证被所有写路径维护又不能表达“删除”事件。数据一旦被物理删除这一行就不存在了增量自然无感知。1.3 删除是周期同步最明显的盲区物理删除对于周期性同步几乎是无解的。全量快照能发现“行数变少了”或“某行不见了”但无法知道它是什么时候被删除的也无法知道删除之前的值是什么。增量同步如果只按时间戳删除操作更是一点影子都看不到。有人会建议“软删除”也就是加一个is_deleted字段删除时把字段置为 1。这能解决一部分问题但并不是所有业务系统都愿意为了数据仓库去改动现有代码。而且软删除会污染查询逻辑所有 SQL 都要额外加WHERE is_deleted 0维护成本很高。binlog 给出了另一个思路MySQL 在主从复制和时间点恢复中本来就记录着每一次 INSERT、UPDATE、DELETE 的完整事件。CDC 把 binlog 当作数据源删除事件不再依赖业务层配合天然就能被捕获到。1.4 主键不变不代表数据没有变周期性同步最容易忽略的一点是很多关键业务字段的变化不体现在“新增行”或“删除行”上而是体现在“同一行的状态翻转”上。典型场景包括用户余额发生变化商品库存扣减和回补订单状态流转配置项被反复修改这些数据的共同点是主键一直存在但值在变。周期同步只给你最后那个值你根本不知道它经历了几次变化、每次变化发生在什么时间。想用这种数据去做“状态停留时长”“变更频率”分析基本不可能。而 binlog 由于记录的是每次变更前和变更后的值你可以用事件流的形态把每一次变化保存下来。这就不是“最终状态表”而是一张可以回放历史的变更事实表。1.5 周期性同步不是错的只是适用边界很清楚我并不是在否定周期同步。对于只需要每日报表、只关心最终状态、对实时性要求不高的场景周期同步仍然是成本最低、链路最简单的方案。它的问题不是说“慢”而是“丢过程”。一旦业务需要回答“为什么变成这样”或者“中间经历了什么”周期同步就会失效。所以在决定是否引入 CDC 之前先要确认你到底需要的是“某个时刻的表”还是“这张表经历过的变化”。2. binlog 为什么能解决“过程缺失”2.1 binlog 是 MySQL 的变更日志binlogBinary Log是 MySQL 服务端生成的一种二进制日志文件主要记录数据库的数据变更。它最早的核心用途有两个主从复制和时间点恢复。主从复制的过程可以简化成主库把每一个数据变更事件写入 binlog从库连接主库并读取这些事件然后在本地重放达到与主库一致的状态。这个机制保证了从库不是简单复制一个结果而是跟着主库的变更顺序在演进。CDC 做的事情本质上和从库“偷学”了同一套机制。它把自己伪装成一个从库去读取 binlog但不直接把变更应用到另一份数据副本而是把这些变更事件转换成消息或数据行发送到下游系统。所以 CDC 不是新发明而是把 MySQL 早就提供的能力用在了数据集成上。2.2 row 格式下 binlog 记录了什么MySQL 的 binlog 有三种记录格式STATEMENT、ROW、MIXED。格式记录内容对 CDC 的意义STATEMENT执行变更的 SQL 语句无法可靠还原每行变化存在不确定性ROW每行变更前后的具体值精确到行能够准确描述变更前后状态MIXED部分语句用 STATEMENT部分自动切 ROW不保证所有事件都精确CDC 不推荐对于 CDC 来说通常要求 binlog 格式改成 ROW。原因很简单ROW 格式下一条 UPDATE 日志会记录这一行在更新前的完整值before image和更新后的完整值after image一条 DELETE 会记录被删除行的所有字段一条 INSERT 会记录新插入行的所有字段。这样下游拿到的不只是一个“变更了”的信号而是可以还原每一次变更的字段级明细。这里有一个容易被忽略的点如果只关心主键STATEMENT 可能也够用但如果要做字段级审计、实时数仓、或者基于变更值做计算那么 before image 和 after image 都是必须的。所以很多 CDC 工具在文档里都会强调要求开启binlog_formatROW binlog_row_imageFULLbinlog_row_imageFULL是为了确保在 update/delete 日志里能拿到完整行数据而不是只记录变化的字段。默认在 MySQL 5.6 之后就是 FULL但也要确认生产环境没有被改过。2.3 从“同步状态”到“同步事件”的视角切换周期性同步解决的是状态复制CDC 解决的是事件捕获。这两者最大的区别在于事件带有顺序、时间和 before/after 信息。当你把 binlog 事件同步到 BigQuery 时得到的不是一个“最新表”而是一段可以重放的变更历史。你可以通过回放这些事件重建任意时刻的表状态也可以只取每个主键最新事件得到当前快照还可以对事件本身做分析统计某类操作发生的频率、时段和条件。这就是为什么说 CDC 不只是“同步得快一点”而是把“数据集成”变成了“数据事件流”。这个视角变化会直接影响你设计下游表结构的方式。如果还是按一张“最终状态表”来做CDC 的价值只发挥了一半只有把变更事件当成独立数据资产才能真正解决过程缺失的问题。3. MySQL CDC 到 BigQuery 的典型链路怎么搭3.1 常见架构选择从 MySQL binlog 到 BigQuery目前社区里比较常见的架构有两种。第一种是使用 Debezium 作为 CDC 组件把 binlog 事件发送到 Kafka再由流处理任务或 Dataflow 模板写入 BigQuery。这套链路的好处是组件成熟、消息有缓冲、可以支撑多个下游系统同时消费坏处是组件多运维成本高。Kafka、Schema Registry、Connect、流处理任务都需要人维护。第二种是使用 Flink CDC 直接读取 MySQL binlog经过计算后通过 BigQuery 连接器写入。这套链路更紧凑减少了 Kafka 这一跳适合已经使用 Flink 的团队。坏处是 Flink 作业的状态管理、checkpoint、重启恢复都需要经验而且 BigQuery 的写入连接器不像 JDBC 那么通用需要根据版本确认配置。如果你所在团队没有专门的流处理平台可以优先考虑云上托管方案比如 Cloud Dataflow 配合 Pub/Sub 或者 Kafka 的模板。这样你不需要自己维护 Flink 集群只需要关注 CDC 配置和 BigQuery 表结构。3.2 用 Flink CDC 的最小流程一个比较典型的最小流程是开启 MySQL binlog设置binlog_formatROW。创建一个有权限的 MySQL 用户至少需要REPLICATION SLAVE、REPLICATION CLIENT、SELECT权限。在 Flink 中定义 MySQL 源表使用mysql-cdcconnector。在 Flink 中定义 BigQuery 目标表使用官方或社区提供的 BigQuery connector。启动一个常驻作业把源表的事件写入 BigQuery。Flink SQL 里的源表定义类似下面这种结构不同版本配置项会有差异落地前以你使用的连接器文档为准CREATE TABLE mysql_orders ( id INT, status STRING, updated_at TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname your-mysql-host, port 3306, username cdc_user, password ******, database-name app_db, table-name orders, scan.startup.mode initial );后面的写入任务可以再把这张源表映射到 BigQuery 对应的表并指定写模式。这里最需要注意的是初始同步阶段Flink CDC 会先做一次历史数据的全量快照然后无缝切换到 binlog 增量。这个切换过程是自动的但如果你给 BigQuery 表定义了主键或唯一约束务必保证历史数据里没有重复行否则初始快照就可能失败。3.3 用 Debezium 加 Kafka 的流程如果选择 Debezium大致的链路是Debezium MySQL Connector 连接 MySQL读取 binlog。把变更事件编码成 JSON 或 Avro发送到 Kafka topic。下游使用 Dataflow、Cloud Function、或者自建 Flink 作业消费这些事件。最终通过 BigQuery 的 Storage Write API 或批量导入接口写入。这里有一个需要想清楚的问题Kafka topic 的消息保存策略是什么。BigQuery 写入如果出现故障Kafka 可以帮你保存一段时间的消息等故障恢复后再继续消费。但如果 topic 的消息过期被清理而下游没有消费完就会出现数据缺口。所以生产环境里Kafka topic 的保留时间一般要大于你预期的故障恢复时间。不过不要为了用 Kafka 而用 Kafka。如果只有一张表要同步并且没有其他下游消费者引入 Kafka 只是增加了运维负担。Flink CDC 或托管模板已经足够。3.4 BigQuery 写入方式流式写入还是批量加载BigQuery 支持两种主要的写入方式批量加载Load Job和流式写入Storage Write API 或传统 streaming insert。批量加载适合离线批量文件成本低性能高但延迟较高。流式写入适合近实时场景数据到达后可以在秒级或分钟级进入可查询状态。CDC 链路通常建议使用流式写入或小批量写入。不建议每条 binlog 事件都立刻写一次 BigQuery这样会产生大量小请求成本和配额都很难受。比较好的做法是攒一批或者让流处理框架按固定间隔、固定条数 flush 一次。Flink 和 Dataflow 里一般都有相关参数类似sink.batch.max-size和sink.batch.interval具体以你使用的 connector 为准。注意不要一上来就把批量数和并发数拉满先用一条样例数据确认输入、输出和日志都正常再逐步调整吞吐。4. 真正落地时会踩的坑4.1 binlog 配置不是默认就能用很多第一次接触 CDC 的团队在 MySQL 里执行了一个SHOW VARIABLES LIKE log_bin发现返回 OFF就直接卡住了。MySQL 默认不一定会开启 binlog即使开启了也不保证格式是 ROW。你需要确认几个关键参数log_bin是否开启 binlog。binlog_format必须是ROW。binlog_row_image建议是FULL。binlog_expire_logs_seconds或expire_logs_daysbinlog 保留多长时间。server_id每个 CDC 连接器需要独立且不与其他复制进程冲突。如果 binlog 没开启或者保留时间太短CDC 任务在中断后可能无法从旧位点继续追只能重新做一次全量快照。这对大表来说代价很高。4.2 server_id 冲突是隐蔽问题CDC 组件读取 binlog 时要扮演一个“从库”角色需要有一个全局唯一的 server_id。如果你在同一个 MySQL 实例上启动了多个使用相同 server_id 的复制进程MySQL 会认为它们是同一个从库event dump 会互相干扰轻则数据重复重则连接被重置。特别是在一个 MySQL 主库上有多个 CDC 任务时一定要给每个任务分配不同的 server_id。这个值可以写在连接配置里也可以由连接器自动生成。如果使用自建 Debezium 或 Flink配置时最好显式指定避免随机冲突。4.3 时区、字段类型和 BigQuery 类型映射MySQL 的TIMESTAMP和DATETIME在 BigQuery 中都需要仔细处理。TIMESTAMP会受 MySQL 会话时区影响DATETIME则没有时区概念。CDC 工具在读取 binlog 时一般会按连接配置的时区进行解析。如果连接串没指定时区解析出来的时间可能和业务实际时间差几个小时。建议在连接配置里显式设置serverTimezone并且让 BigQuery 侧统一使用 UTC 或某个固定时区存数。不要在源库、连接器、目标表三层各用各的时区否则排查时间问题时非常痛苦。字段类型方面比较容易踩坑的有MySQL 的tinyint(1)在部分驱动里会映射成Boolean如果业务里有多个tinyint要逐一确认。unsigned类型在 CDC 事件里的值可能超过目标类型上限需要注意映射到INT64还是NUMERIC。DECIMAL精度要和 BigQuery 的NUMERIC或BIGNUMERIC匹配否则会有舍入问题。4.4 删除事件不该直接物理删除BigQuery 支持 DML 操作理论上可以执行DELETE语句。但 CDC 场景下如果每次上游删除都转成 BigQuery 的 DELETE成本很高而且和流式写入的模型不匹配。更常见的做法是在目标表增加一个_is_deleted标志字段。上游 DELETE 事件到达时不实际删除目标行而是把_is_deleted设为 true同时把_event_timestamp更新为删除时间。查询时默认过滤_is_deleted false需要审计时再查全量。这样做的好处是保留审计历史也避免 BigQuery 的频繁 DML。缺点是目标表会越来越大需要根据保留策略做分区和清理。4.5 schema 变更处理MySQL 表结构不会一直不变。上线新功能时很可能会加字段、改字段长度、甚至改字段名。这会让 CDC 事件里的 JSON 或字段数量发生变化而 BigQuery 目标表结构如果不同步更新写入就会失败。建议提前规划上线前通知数据团队同步修改 BigQuery 目标表结构。使用支持 schema evolution 的 CDC 工具和连接器比如某些 Debezium 配置可以自动加字段。如果不支持自动变更至少要有一个监控机制发现schema mismatch报错后及时处理。最怕的是有人改了 MySQL 表结构但 CDC 任务还在运行事件里多出来的字段导致写入报错然后整个任务积压。等发现时binlog 可能已经推进了很远的位点。4.6 数据延迟、积压和重放CDC 链路是流式的上游 binlog 不会等下游处理完。如果 BigQuery 写入速度跟不上或者 Kafka 消费变慢事件就会积压。积压时间过长可能出现内存溢出、checkpoint 失败、任务重启。这里需要区分两种“重复”消费重试导致的重复和源库本身在同一时刻对同一行进行了多次变更。后者是正常业务事件前者是系统问题。大多数 CDC 管道的投递语义是“至少一次”也就是极端情况下会重复消费。你需要在目标表设计上去重或者使用 upsert 语义保证最终一致性。5. 排查链路和长期运维建议5.1 当发现数据不对不要先调参数很多团队看到 BigQuery 数据比 MySQL 少第一反应是加大 Flink 的并发或调小批处理间隔。但大多数时候问题不在吞吐。我建议按这个顺序排查先看 CDC 任务状态是否正常运行最近有没有 checkpoint 失败当前消费位点和延迟指标是多少。再看 MySQL 侧binlog 是否还在连接用户是否有对应权限server_id 是否冲突。再看消息链路事件是否到达了 Kafka topic 或流计算源表有没有解析异常。再看 BigQuery 写入任务最近有没有报错是不是结构不匹配是不是配额超限。最后做数据校验对比源表和目标表的主键数量、关键字段抽样、最近更新时间。顺序的逻辑是先确认链路没有断再确认数据有没有被正确解析最后才考虑是不是写入层的问题。如果一上来就调整参数可能掩盖了真正的问题甚至让任务在新的参数下产生更多重复数据。5.2 运维该盯哪些指标CDC 任务不是搭完就结束的。长期维护时至少要关注这些指标lag当前消费的 binlog 位点离 MySQL 最新位点有多远。checkpoint周期和失败次数流处理里的数据一致性依赖 checkpoint。dead letter queue数量解析或类型转换失败的事件是否在堆积。BigQuery 写入错误配额、超时、schema mismatch。binlog 保留时间确认你没快照完的 binlog 不会被提前清理。如果这些指标有告警数据问题通常能在业务反馈之前被发现。5.3 先跑通、再优化、最后工程化一个稳健的推进路径是这样先选一张关键小表配置好 binlog 和 CDC 作业把数据同步到 BigQuery 的临时表。验证源表和目标表的数据一致性确认更新、插入、删除三类事件都正确。再扩展到需要实时数据的核心大表观察延迟和资源占用。最后再把备份、告警、重放、schema 变更流程补齐。不要第一天就想把所有表都接入 CDC。一次接太多表一旦发生 schema 变更或者 binlog 积压排查范围会非常大反而拖慢整体进度。6. 什么时候应该用 CDC什么时候继续用周期同步6.1 三问模型在为一个 MySQL 到 BigQuery 的同步项目做技术选型时我会先问三个问题下游是否需要看到“变更中间态”或历史变化数据时效是否需要达到分钟级甚至秒级团队是否有能力运维流式链路如果三个答案都是否那么周期同步仍然是一个效率很高的选择。你不需要因为 CDC 是热门概念就强行引入。如果三个答案里至少有两个是“是”那么 CDC 大概率是值得投资的方案但也不一定要覆盖所有表。6.2 从成本、运维、数据质量三个维度看维度周期同步CDC成本定时任务成本低需要额外的流计算、消息队列、存储成本运维简单失败重跑即可需要关注位点、积压、schema 变更、幂等数据质量只能给出最终状态容易丢过程保留完整变更事件但需要处理重复和类型映射时效取决于同步周期可以做到分钟级甚至秒级适用场景日报、周报、最终状态不变审计、历史回放、实时看板、实时告警这张表不是想说 CDC 一定更好而是说两者回答的是不同问题。6.3 最小化建议如果你还在犹豫可以先只对最核心的业务表启用 CDC比如订单、用户余额、库存。其他那些很少变化的维度表保持周期同步完全足够。这样可以先用最低成本验证 CDC 链路是否稳定也更容易控制风险。同时把原有的周期同步任务留作兜底至少在初期不要马上删掉。CDC 链路刚上线时可能还有类型映射、时区、重复等问题有周期同步做交叉比对能更快发现问题。回到最开始的场景。如果你现在正为 BigQuery 和 MySQL 数据对不上而头疼先不要急着把同步频率从每天一次改成每小时一次。频率再高也只是一种更细粒度的快照仍然无法回答“这行数据经历了哪些变化”这个问题。真正需要你先想清楚的是下游要的到底是某个时刻的表还是这张表经历过的每一次变化。这个判断清楚了binlog、CDC、BigQuery 的组合才不是一个跟风的架构选型而是真正解决问题的工具箱。
返回列表