
MySQL 到 BigQuery 的数据同步最容易被低估的问题就是时间窗口。无论是定时导出还是按updated_at增量拉取本质上都属于 periodic syncs。它们的共同点是数据库里的变化并不会等待调度任务开始也不会按周期整齐地落入边界。一次删除、一条字段被改回旧值、一张表在夜间被大批量 UPDATE 后又改回来这些事件都可能发生在两批同步任务的间隙最终 BigQuery 里的数据既不是源表的真实状态也不是任何历史时刻的真实状态。CDCChange Data Capture通过读取 MySQL binlog把每一条数据变更作为事件流送到 BigQuery正好从机制上补上了这个缺口。这篇文章围绕周期同步会漏什么、binlog 为什么能避免漏、落地时要注意什么展开适合正在设计数据管道、给数仓接增量数据或者被批量任务数据不一致问题困扰的开发者与数据工程师。1. 周期同步在 MySQL 到 BigQuery 场景下到底漏了什么1.1 常见的三种周期同步写法先看最常用的三种同步方式它们并不只是实现细节不同能捕获的数据变化粒度也完全不同。第一种是全量导出覆盖。直接把 MySQL 表导出成文件或通过 SQL 拉取写入 BigQuery 临时表再覆盖目标表。这种方式能保证目标表最终状态一致但同步窗口很长且 BigQuery 做覆盖时下游可能读到一半数据。数据量一旦上亿这个方案基本不可持续。第二种是按自增 ID 增量拉取。记录max(id)每次只拉大于该 ID 的行像这样SELECT * FROM orders WHERE id :last_max_id ORDER BY id;这个方案只能捕获新增数据。业务表一旦发生 UPDATE主键 ID 不变增量 SQL 永远拉不到这一行。DELETE 更不会出现在结果里。第三种是按更新时间戳增量拉取SELECT * FROM orders WHERE updated_at :last_sync_ts;这是目前最常见的周期同步方案前提是业务表有updated_at字段并且所有写入路径都正确更新这个字段。实际项目里这个前提经常被破坏某些批量导入脚本没有更新updated_at某些框架写入时没有映射该字段于是出现数据明明变了增量 SQL 却查不到的问题。1.2 周期同步一定会错过的几类变更物理删除是最典型的一类。DELETE 之后这条记录不再存在于表中任何基于当前表状态的 SELECT 都无法发现它曾经存在过。全量对拍能发现问题但只能事后补救而且对拍本身在大表上成本极高。没有更新时间字段或者更新时没有写入时间戳也是一类。订单表如果通过第三方系统直接改库或者 DBA 手工执行 UPDATE 时没有维护updated_at那么增量边界从一开始就是错的。同周期内状态回跳同样会被掩盖。假设订单在 00:00:10 从pending改为paid00:00:20 又改回pending。周期任务在 01:00 运行拉到的最终状态还是pending。从业务角度看中间那次paid状态也曾经是真实数据但周期同步完全感知不到。高频更新更不用说。一张促销表每秒更新几千行周期任务每隔 5 分钟拉一次单行在周期内被反复更新后最终拉到的只是最后一次值中间所有取值全部丢失。1.3 为什么不是多跑几次就能解决周期同步的失败模式是逻辑性漏数据不是漏跑任务。把调度频率从小时改成分钟只是缩小时间窗口并没有改变读取当前表状态的本质。一张表在周期内发生了 100 次更新周期同步只能看到最后一行binlog 能看到 100 个事件并且每个事件都保留前镜像和后镜像。这就是原理层面的差异。周期同步试图通过查询结果反推变化而 binlog 是 MySQL 自己记录的写操作流水账。流水账不会因为业务表没有updated_at就缺页也不会因为 DELETE 后记录消失就抹去历史。1.4 三种方案的能力对比维度全量快照增量字段轮询binlog CDC删除事件全量对拍后才发现通常无法发现每条 DELETE 都有对应事件更新历史只有最后状态只有最后一次变更每次 UPDATE 都有前镜像和后镜像对业务表要求无必须有updated_at等字段无binlog 与业务表结构独立实时性取决于调度周期取决于调度周期秒级到分钟级可配置对源库压力大全表扫描代价高中等取决于索引较小读取日志而不是反复扫描表从这张表能看出周期同步不是慢而是漏。CDC 的价值不是让同步更快而是让变化过程本身可见。2. binlog 为什么能捕捉每一次变化CDC 的原理2.1 binlog 是什么binlog 是 MySQL 的二进制日志记录所有改变数据库内容的操作包括 INSERT、UPDATE、DELETE以及部分 DDL。MySQL 主从复制、崩溃恢复、数据恢复都依赖它。可以理解为 MySQL 把每一次写操作按顺序写到一本流水账上。binlog 并不是默认可用的。MySQL 5.7 中log_bin默认关闭8.0 默认开启但不同发行版和云厂商的默认值可能不同落地前必须先确认。如果 binlog 没有开启后续所有 CDC 方案都无从谈起。2.2 ROW 格式给 CDC 提供了什么binlog 有三种格式STATEMENT、ROW、MIXED。STATEMENT 格式记录的是 SQL 语句本身例如UPDATE orders SET statuspaid WHERE id1001;。这种格式日志量小但无法可靠还原每一行在语句执行前后的具体值。MIXED 格式是两者的混合MySQL 会根据语句类型自动选择但对于 CDC 场景依然不够稳定。CDC 要求使用 ROW 格式。ROW 格式下binlog 直接记录行的变化包括字段级的前镜像和后镜像。具体来说INSERT 事件包含插入后的完整行数据。UPDATE 事件包含变更前的整行数据和变更后的整行数据。DELETE 事件包含删除前的整行数据。这意味着 CDC 消费者不仅能知道某张表发生了变化还能拿到 哪一行的哪个字段从什么值变成什么值。2.3 CDC 连接器如何消费 binlogDebezium、Flink CDC 这类工具在原理上会伪装成 MySQL 从库。它们通过 MySQL 的复制协议从主库拉取 binlog并把 binlog 里的二进制事件解析成结构化的 JSON 变更事件。连接器需要记录自己的消费位点。传统方式是记录 binlog 文件名加偏移量例如mysql-bin.000023的position 45123。更可靠的方式是使用 GTID即全局事务标识符。GTID 能唯一标识每个事务即使 binlog 文件被清理只要 MySQL 实例保留了完整的事务历史连接器也能定位到正确的起点。CDC 连接器通常具备先快照再增量的能力。首次启动时它会先读取一次源表全量数据记录当时的 binlog 位点之后继续从该位点消费增量从而保证从启动那一刻起不遗漏后续变更。2.4 从 binlog 到 BigQuery 的完整链路一个常见的生产架构是MySQL master - binlog - CDC Connector (Debezium / Flink CDC) - Kafka Topic - 流处理或写入程序 - BigQuery Storage Write API / Load Job - BigQuery Table也可以简化为MySQL master - Flink CDC - BigQuery 目标表无论采用哪种架构核心都是从日志读取变化而不是定时查询表。这也决定了后面的环境准备、配置、验证和排错方式。3. 前期准备MySQL、BigQuery 和权限一项都不能省3.1 版本与前置条件在配置 CDC 之前先确认环境是否满足基本条件组件要求说明MySQL5.7 或 8.0开启 binlog5.7 建议显式开启8.0 确认默认配置BigQuery数据集、目标表、服务账号建议单独建服务账号避免共用管理员账号CDC 工具Debezium 或 Flink CDC版本需要与 MySQL 和 Kafka 版本匹配网络源库与数仓侧连通私网优先公网场景需要做好传输加密如果源 MySQL 是云数据库还需要查看云厂商是否允许开启 binlog 保留策略、是否开放复制账号权限。有些托管数据库默认不开放REPLICATION SLAVE这是接入 CDC 前最容易发现的阻塞点。注意开启 binlog 并切换为 ROW 格式后binlog 日志量通常会变大磁盘占用和复制延迟都会上升。生产环境切换前需要评估磁盘余量。3.2 修改 MySQL 配置下面是一份最小可用的 MySQL CDC 配置示例[mysqld] server_id 1001 log_bin /var/log/mysql/mysql-bin.log binlog_format ROW binlog_row_image FULL expire_logs_days 7 # MySQL 8.0 可用以下参数控制 binlog 保留时长 # binlog_expire_logs_seconds 604800每个参数的作用server_idMySQL 实例在复制拓扑中的唯一标识。CDC 客户端也会占用一个 server-id不能与主从库中其他节点重复。log_bin开启 binlog并指定日志文件路径。binlog_formatROW让 binlog 记录行级变更。CDC 必须使用 ROW 格式。binlog_row_imageFULL让 UPDATE 事件包含整行前镜像和后镜像。如果设置为 MINIMALbinlog 只包含被修改的字段和主键CDC 拿不到完整旧行和新行。expire_logs_days控制 binlog 文件保留天数。保留太短CDC 位点落后时可能追不上保留太长磁盘占用过大。常见建议是 3 到 7 天具体要结合源库写入量和磁盘容量调整。修改配置后需要重启 MySQL。重启前确认max_allowed_packet等参数不会限制大事务的 binlog 传输。3.3 创建 MySQL CDC 账号建议为 CDC 单独创建一个账号避免使用 rootCREATE USER cdc_user% IDENTIFIED BY strong_password; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES;三个权限的含义SELECT用于 CDC 工具首次启动时的全量快照以及读取表结构信息。REPLICATION SLAVE允许该账号通过复制协议读取 binlog这是 CDC 的核心权限。REPLICATION CLIENT允许执行SHOW MASTER STATUS、SHOW BINARY LOG STATUS等命令用于确认位点信息。不要把ALL PRIVILEGES都授出去。CDC 账号只需要读取能力不需要写源库。3.4 BigQuery 侧准备BigQuery 侧需要准备数据集、目标表和服务账号。在 Google Cloud Console 中先创建数据集例如analytics。目标表建议在接入 CDC 之前就定义好字段类型尽量与 MySQL 类型对应。如果后续依赖 BigQuery 自动加列容易遇到 schema 不一致导致写入失败。服务账号需要授予 BigQuery Data Editor 或更细粒度的角色。把服务账号的 JSON 密钥下载到写入服务所在机器并通过环境变量GOOGLE_APPLICATION_CREDENTIALS指向密钥文件。BigQuery 是列式存储目标表 schema 在写入前就要对齐。CDC 事件字段如果比目标表多需要做过滤如果少目标表多出的列会使用默认值或 NULL。4. 最小落地链路Debezium 捕获 binlog程序写入 BigQuery4.1 两种常用的技术选型常见方案有两种方案链路适合场景Debezium KafkaMySQL - Debezium - Kafka - 写入程序 - BigQuery已有 Kafka 基础设施需要多消费方Flink CDCMySQL - Flink CDC - BigQuery Sink团队熟悉 Flink希望用 SQL 处理流下面以 Debezium Kafka Python 消费者为例把链路拆开看。这样更容易理解每个环节的职责。Flink CDC 只是把 Debezium 和流处理合并到一个框架里原理一致。4.2 Debezium connector 的配置Debezium 通过 Kafka Connect 运行一个典型配置如下{ name: mysql-orders-connector, config: { connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: 10.0.0.10, database.port: 3306, database.user: cdc_user, database.password: xxxx, database.server.id: 5400, database.include.list: ecommerce, table.include.list: ecommerce.orders, database.history.kafka.bootstrap.servers: kafka:9092, database.history.kafka.topic: schema-changes.ecommerce, topic.prefix: mysql, include.schema.changes: true } }关键参数database.server.idDebezium 会占用一个 server-id。它必须与 MySQL 现有主从库、其他 CDC 实例的 server-id 不冲突否则连接会被 MySQL 拒绝。database.include.list/table.include.list限定监听的库表。只同步需要的表能显著减少 binlog 解析压力。database.history.kafka.topicDebezium 用这个 topic 记录表结构历史。binlog 里的旧事件在解析时可能依赖历史 schema因此这个 topic 不能随意删除。topic.prefix生成 Kafka topic 名称的前缀。最终 topic 名称一般是{topic.prefix}.{database}.{table}。4.3 变更事件长什么样Debezium 输出的变更事件是一段 JSON核心结构如下{ before: { id: 1001, status: pending }, after: { id: 1001, status: paid }, source: { db: ecommerce, table: orders, server_id: 1001, ts_ms: 1719900000123 }, op: u }op字段表示操作类型op 值含义事件内容cINSERT只有afteruUPDATE有before和afterdDELETE只有beforer快照读取类似 INSERTafter为快照行注意DELETE 事件没有after。写入 BigQuery 时如果目标表要反映删除必须自己定义删除策略比如写入一条带删除标记的记录或者通过主键 MERGE 删除目标行。4.4 写入 BigQuery 的示例程序下面是一个最小 Python 消费者示例从 Kafka 读取 MySQL 变更事件批量写入 BigQueryimport json from google.cloud import bigquery from kafka import KafkaConsumer PROJECT my-project DATASET analytics TABLE orders client bigquery.Client(projectPROJECT) table_ref client.get_table(f{PROJECT}.{DATASET}.{TABLE}) def process_event(msg): payload json.loads(msg.value()) op payload.get(op) if op in (c, r): return payload[after] if op u: return payload[after] if op d: before payload[before] before[_is_deleted] True return before return None consumer KafkaConsumer( mysql.ecommerce.orders, bootstrap_serverskafka:9092, group_idbigquery-sync, auto_offset_resetlatest, enable_auto_commitFalse, ) rows [] batch_size 500 for message in consumer: row process_event(message) if row is not None: rows.append(row) if len(rows) batch_size: errors client.insert_rows_json(table_ref, rows) if not errors: consumer.commit() rows [] else: print(errors)这个示例说明的是思路不是完整生产代码。insert_rows_json适合小规模验证生产环境更推荐使用 BigQuery Storage Write API并配合监控、重试和死信队列。enable_auto_commitFalse是为了避免消息未成功写入就提交位点减少丢失风险但代价是重复消费因此目标表必须容忍重复。4.5 如果团队已经用 Flink可以考虑 Flink CDCFlink CDC 可以把上面的链路压缩成一个 SQL 和一套连接器。用 Flink SQL 创建 MySQL CDC 源表CREATE TABLE mysql_orders ( id INT, user_id INT, amount DECIMAL(10, 2), status STRING, updated_at TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 10.0.0.10, port 3306, username cdc_user, password xxxx, database-name ecommerce, table-name orders, server-id 5400-5406, scan.incremental.snapshot.enabled true );scan.incremental.snapshot.enabled在较新版本默认开启。它让 Flink CDC 以分片方式并行快照大表不需要像旧版本那样先对全表加锁再读取对大表更友好。源表创建后可以再创建 BigQuery Sink 表通过INSERT INTO完成同步。具体 Sink 类名和参数取决于连接器版本落地前要以当前使用的 Flink 和连接器文档为准。5. 怎么验证 binlog 同步没有漏数据5.1 先确认 binlog 真的开了进入 MySQL 命令行执行SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;预期结果中log_bin为ONbinlog_format为ROWbinlog_row_image为FULL。还可以执行SHOW BINARY LOG STATUS;如果输出包含当前 binlog 文件名和 position说明 binlog 文件正在正常写入。5.2 验证 Kafka 收到了哪些变更先用 Kafka 自带的控制台消费命令观察 MySQL 变更是否进入 topickafka-console-consumer.sh \ --bootstrap-server kafka:9092 \ --topic mysql.ecommerce.orders \ --from-beginning然后在 MySQL 中分别执行一次 UPDATE 和一次 DELETEUPDATE orders SET status paid WHERE id 1001; DELETE FROM orders WHERE id 1002;正常情况下消费端会看到op为u和d的两条事件。这一步直接验证了周期同步最难做到的能力删除和更新都能被捕获。5.3 验证 BigQuery 目标表观察 BigQuery 目标表是否有新数据写入。可以通过控制台查询也可以执行SELECT COUNT(*) FROM my-project.analytics.orders; SELECT MAX(updated_at) FROM my-project.analytics.orders;必须注意一个容易误判的地方BigQuery 目标表不会因为收到了 DELETE 事件就自动删除对应行。如果写入程序只是把after或before以追加方式写入删除事件只会变成一行带标记的数据。要真实反映删除目标表需要按主键做 MERGE或者通过分区覆盖实现。验证时先明确自己的目标表语义是追加明细还是镜像源表。5.4 延迟监控指标从 binlog 到 BigQuery 的同步不是一次性的必须持续监控。常见指标包括指标含义告警建议Kafka consumer lag消费程序落后的消息数持续增长则告警Debezium 位点与当前 binlog 的文件间隔连接器是否在追赶超过 binlog 保留期则高风险端到端延迟事件写入 MySQL 到进入 BigQuery 的时间差根据业务要求设置阈值BigQuery 写入错误率schema 不匹配等写入失败立即告警把位点落后和consumer lag 持续增长作为关键告警能提前发现大事务、网络抖动或消费程序故障。6. 数据到达 BigQuery 后模式映射、DDL 和幂等才是真正的坑6.1 MySQL 与 BigQuery 类型映射字段类型映射是 CDC 链路里最容易踩坑的部分。下面是常见映射关系MySQL 类型BigQuery 类型注意事项INT / INTEGERINT64无符号 INT 可能超过 INT64 有符号范围BIGINTINT64超过 2^63-1 的数据要改用 NUMERIC 或 STRINGDECIMAL(p, s)NUMERIC / BIGNUMERIC金额字段不要用 FLOAT精度会丢失DATETIMEDATETIME无时区语义按原值写入TIMESTAMPTIMESTAMP建议统一按 UTC 存储VARCHAR / TEXTSTRING长度和编码要注意JSONJSONBigQuery 需要字段模式为 JSON 或先转成 STRINGTINYINTINT64 / BOOL看业务语义确定最容易出问题的是 DECIMAL。MySQL 中的DECIMAL(10, 2)如果映射成 BigQuery 的 FLOAT640.1 这样的值可能出现精度误差。正确做法是映射为 NUMERIC。TIMESTAMP 也容易出问题。MySQL 的TIMESTAMP有会话时区概念CDC 事件里的ts_ms可能是 UTC 时间而业务字段本身可能是本地时间。建议在写入端统一规范避免目标表同一列混入不同时区的数据。6.2 DDL 变更会打断 CDC当 MySQL 表结构变化时CDC 链路会面临两个层面的问题。第一Debezium 需要依赖database.history.kafka.topic中的 schema 历史来解析 binlog 里的旧事件。如果这个 topic 被删除或清理连接器可能无法反序列化旧的 binlog 事件。第二BigQuery 目标表的 schema 不会自动跟随 MySQL DDL 变化。MySQL 加了一列CDC 事件里出现了新字段但 BigQuery 目标表没有这一列写入就会报错。处理建议是把 DDL 纳入变更流程先审查 MySQL DDL 对同步链路的影响。先在 BigQuery 目标表补充或调整 schema。再在 MySQL 执行 ALTER TABLE。同步完成后核对事件是否正常。对于大表的 ALTER TABLE还可能导致源库锁表和复制延迟。生产环境做主从切换时要评估 DDL 对 binlog 位点的影响。注意不要依赖 BigQuery 自动加列来处理所有 DDL 变更。自动加列在不同版本和连接器里行为不一致且不能处理列重命名、删除、类型变更等复杂操作。6.3 至少一次语义下重复是正常的binlog CDC 链路通常提供 at-least-once 语义。网络闪断、消费程序重启、位点提交失败都可能导致同一事件被重复消费。因此目标表必须能接受重复。常见做法按主键去重写入前先判断目标表是否已有该主键。使用 BigQuery MERGE按主键更新目标行。在记录中增加事件版本字段如event_ts_ms或 GTID写入时取较新的事件。下面是 BigQuery MERGE 的简化思路MERGE my-project.analytics.orders AS t USING changes AS s ON t.id s.id WHEN MATCHED THEN UPDATE SET status s.status, amount s.amount WHEN NOT MATCHED THEN INSERT (id, user_id, amount, status, updated_at) VALUES (s.id, s.user_id, s.amount, s.status, s.updated_at);MERGE 在处理删除事件时还可以加一个WHEN MATCHED AND s._is_deleted TRUE THEN DELETE分支。但 MERGE 的成本比流式追加高适合对一致性要求高、更新频率可控的场景。如果表更新量极大需要考虑分区覆盖、冷热分离等方案。6.4 乱序事件怎么处理同一个主键的多条变更在 Kafka 中如果分布到不同分区消费程序收到的顺序可能和源库事务提交顺序不一致。比如先提交了statuspaid后提交了statuscancelled乱序可能导致目标表最终停在paid。处理方式Kafka Topic 按主键 hash 分区保证同一主键路由到同一分区。写入端使用 binlog 里的ts_ms或 GTID 做排序只接受更新的事件。如果业务允许短暂延迟可以在写入端做窗口缓冲按主键排序后批量提交。如果源表存在删主键后重新插入同一主键的场景还需要区分删除后插入和旧 UPDATE 后到否则可能出现旧数据覆盖新数据的现象。这种情况下GTID 或事务 ID 是更可靠的顺序依据。7. 常见问题排查从现象倒推 binlog 链路故障7.1 现象连接器启动时报权限不足或无法读取 binlog可能原因MySQL 账号缺少REPLICATION SLAVE权限。连接器配置的 server-id 与现有从库冲突。binlog 未开启或者binlog_format不是 ROW。排查命令SHOW VARIABLES LIKE binlog_format; SHOW GRANTS FOR cdc_user%; SHOW PROCESSLIST;处理方式核对 MySQL 配置和账号权限修改后重启连接器。server-id 冲突通常会在 MySQL 错误日志里看到A slave with the same server_uuid/server_id as this slave has connected to the master之类的信息。7.2 现象任务运行一段时间后Kafka 里有历史事件但新事件迟迟不来可能原因Kafka Connect 或连接器进程挂掉后位点没有正确恢复。MySQL 实例重启导致 binlog 文件名变化连接器找不到旧位点对应的文件。table.include.list配置了大小写敏感的表名实际表名大小写不一致。排查方式kafka-consumer-groups.sh --bootstrap-server kafka:9092 --describe --group bigquery-sync重点看CURRENT-OFFSET、LOG-END-OFFSET和LAG。如果 consumer lag 为 0 但新数据没进来检查连接器日志里 binlog offset 是否还在推进。必要时做一次 重新快照 增量 的初始化。7.3 现象BigQuery 写入报错字段不存在或类型不匹配可能原因MySQL DDL 新增了列BigQuery schema 没有同步。DECIMAL 字段映射成了 FLOAT64导致精度丢失或写入失败。MySQL JSON 字段映射到了 BigQuery STRING但事件里是 JSON 对象。排查方式SELECT column_name, data_type FROM my-project.analytics.INFORMATION_SCHEMA.COLUMNS WHERE table_name orders;处理方式定位是哪一列不匹配先同步 schema再重放失败事件。不要直接丢弃报错事件否则会在对账时发现数据缺口。7.4 现象同步延迟持续增长可能原因源库执行了大事务例如一次 UPDATE 超过十万行binlog 事件量巨大。消费程序单线程写入 BigQuery写入速度跟不上源库变更速度。网络带宽不足或者 BigQuery 写入配额受限。处理方式在源库侧避免一次性更新超大范围拆成小事务。写入端改用批量并行写并启用 Storage Write API。增加监控观察 binlog 保留时间是否充足。如果消费端位点落后太远而 binlog 文件已经过期可能需要重新快照。