ARTICLE DETAIL

资讯详情

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

MySQL binlog 到 BigQuery 的 CDC 实时同步:从原理到实践

MySQL binlog 到 BigQuery 的 CDC 实时同步:从原理到实践 这次我们不聊抽象概念专门看一个数据链路里的实际问题MySQL 数据同步到 BigQuery很多团队第一版都是定时批量同步跑了一段时间后会发现报表对不上、明细缺行、删除和更新根本没同步过去。问题不在写 SQL 的人而在同步思路本身。要解决这类问题需要把 MySQL 的 binlog 当成数据源用 CDCChange Data Capture的方式把数据库的每一次变更准实时地送进 BigQuery。本文会先拆定期同步漏数据的具体场景再讲 binlog 的基础格式然后给出一套可落地的 CDC 架构包含环境准备、部署步骤、接口/批量任务说明、性能观察和常见问题排查。适合已经在用 BigQuery 做数仓、用 MySQL 做业务库现在想把同步链路升级成实时或准实时的数据工程师。1. 核心能力速览能力项说明同步方式基于 MySQL binlog 的增量变更捕获实时推送至 BigQuery解决的问题删除/更新/物理清理漏同步、批量窗口内数据不一致、报表延迟高数据源要求MySQL 5.7/8.0 等版本需开启 binlog建议使用 ROW 格式目标端BigQuery 表天然适合 Append 事件表和按主键 Upsert 的镜像表常用方案Canal、Debezium Kafka、Flink CDC、Dataflow 等可按团队技术栈选择部署方式Docker、云托管连接器或自建服务均可对接是否支持 API支持CDC 组件通常提供 HTTP/REST 管理接口或通过消息队列对外输出是否支持批量任务支持存量全量导入 增量实时同步可组合执行延迟水平秒级到分钟级取决于组件链和消费端写入能力适合场景实时数仓、变更追踪、审计分析、主数据同步、缓存/搜索索引联动这里先给结论如果你的表需要“源端改了目标端马上能看到”或者需要完整保留更新和删除动作定时同步不是修修补补能解决的应该直接切换到 binlog 驱动的 CDC。2. 为什么定期同步会漏数据定期同步最常见的形式是每小时或每天跑一次任务把 MySQL 里符合条件的行全量拉出来然后覆盖写进 BigQuery。这种方案在数据量小、表结构稳定、没有删除操作时能跑但一旦业务上出现下面几种情况问题就暴露了。2.1 同步窗口内的中间变化被覆盖假设昨天 23:00 跑了一次全量同步今天 23:00 又要跑第二次。中间 24 小时里一条订单记录被创建、修改了 5 次最后又取消。定时任务最终看到的只是“取消”后的最终值中间 5 次状态流转全部丢失。对于订单、审核、物流这类需要过程分析的场景这个损失是致命的。2.2 删除操作无法被完整感知全量同步的目标表如果是“先清空再写入”或者“按主键覆盖”你会发现源表物理删除的行在目标表里可能仍然存在。因为全量同步看到的只是当前快照它不知道哪些行已经消失。即使做全量比对也只能在下次对账时发现差异无法知道这条记录是什么时间、被谁删除的。2.3 定时批次本身存在延迟窗口每天一次同步意味着 BigQuery 里的数据天然滞后 24 小时。就算改成每小时一次仍然有“批次边界”问题批处理任务执行到一半时源库还在发生新变更这次任务拉到的数据可能不完整下次任务又可能因为主键重复或时间戳边界和上次重叠造成重复或缺失。2.4 失败重跑可能造成重复或脏数据定时任务经常面临“凌晨跑失败了早上手动补跑”的情况。补跑如果只按时间窗口重拉已经写入的数据和新拉到的数据之间没有天然的去重机制。需要在目标表上维护一个“最新更新时间”的判断逻辑这会让 SQL 越来越复杂逐步变成一个不可维护的状态。2.5 DDL 变更会让批量任务直接崩掉业务表加了一个字段或者改了字段类型定时任务的 SELECT * 可能立刻失败或者字段错位写入 BigQuery。更麻烦的是批量任务通常不会记录表结构版本一旦源端 DDL 变化后续任务很容易持续失败。这一节的核心结论定时同步本质上是在对比“两个时刻的快照”它天然丢失过程、丢失删除、延迟明显也扛不住 DDL 变化。要拿到连续、完整、可回放的变更流必须回到 MySQL 内部已经存在的 binlog。3. binlog 是什么为什么它适合做同步依据binlog 是 MySQL 的二进制日志记录的是数据库层面的所有数据变更逻辑。MySQL 主从复制依赖它数据恢复也依赖它而我们做 CDC 同样是消费它。3.1 binlog 的三种格式格式内容对 CDC 的意义STATEMENT记录执行过的 SQL 语句不适合 CDC因为无法精确还原数据行变化ROW记录每行变更前和变更后的值最适合 CDC删除、更新、插入都能完整还原MIXED由 MySQL 根据情况自动选择建议不用在 CDC 链路中可能出现部分语句格式不一致从材料看绝大多数生产环境在做 CDC 时都会把 binlog_format 设为 ROW原因很简单一条 UPDATE 语句在 ROW 格式下会生成多个事件每个事件包含被修改行的 before 和 after 数据消费者不需要解析原始 SQL直接拿到精确值。3.2 需要关注 binlog 中的事件类型CDC 链路真正关心的主要事件包括INSERT_ROWS_EVENT新插入的行数据。UPDATE_ROWS_EVENT更新前后的行数据。DELETE_ROWS_EVENT被删除的行数据。TABLE_MAP_EVENT事件对应的表结构信息。QUERY_EVENT / DDL_EVENTDDL 变更记录用于感知表结构变化。GTID_EVENT如果开启 GTID可以拿到全局事务标识方便精确判断消费位点。3.3 为什么 binlog 比时间戳字段更可靠很多人会问如果表里本来就有 updated_at 字段能否用“查询 updated_at 上次同步时间”代替 CDC这种做法有两个硬伤。第一物理删除的行不会保留 updated_at删除事件直接丢失第二如果业务代码更新时没写 updated_at或者使用了数据库直接修改、存储过程批量更新时间戳字段可能完全不可信。binlog 是 MySQL 自己生成的不依赖业务逻辑是否规范所以可靠性更高。3.4 GTID 和位点binlog 消费时需要记住消费位置。MySQL 提供了两类位置信息传统的文件名 偏移量以及 GTID 集合。生产级 CDC 工具通常同时支持这两种推荐开启 GTID。它在故障恢复、链路切换时更稳定不容易因为 binlog 文件清理或主从切换导致消费位置失效。4. MySQL 到 BigQuery 的 CDC 整体架构CDC 链路的基本流程是四个阶段MySQL binlog - 捕获组件 - 传输组件 - BigQuery 写入组件捕获阶段读取 binlog 并解析成结构化事件传输阶段负责把事件送到下游可能经过消息队列写入阶段在 BigQuery 里完成追加或合并。4.1 常用方案对比方案组件特点适合团队Canal 自定义写入Canal 解析 binlog业务服务消费轻量Java/Scala 技术栈友好需要自己写下游写入已有 Java 服务不想引入 KafkaDebezium KafkaDebezium Connector Kafka Connect生态成熟事件格式标准可接入多个消费端已有 Kafka 基础设施Flink CDCFlink SQL / YAML 作业支持 SQL 直接定义同步批流一体方便回放和清洗倾向用 Flink 做实时计算Dataflow Pub/Sub使用 Pub/Sub 接收变更Dataflow 写 BigQuery云上托管运维少适合 Google Cloud 深度用户目标就是全托管选型时不用纠结“哪个最好”而要看团队已有的运行时。如果用 Kafka接 Debezium 最顺如果已经在用 Flink 做实时计算Flink CDC 能少一个转发环节如果没有实时计算平台Canal 或轻量脚本更容易落地。4.2 BigQuery 端模型设计CDC 到 BigQuery 后目标表有两种常见模型。第一种是事件流水表Append 模型每个变更事件一行包含变更前数据、变更后数据、操作类型、binlog 时间戳。适合审计、过程分析和离线回溯。{ op: update, source_table: orders, primary_key_id: 1001, before: {status: pending}, after: {status: paid}, event_time: 2025-01-15 10:30:00 }第二种是镜像表Upsert 模型BigQuery 里保留业务表当前最新状态每次变更事件按主键覆盖或合并。适合 BI 报表直接查询延迟低不需要每次跑全量。实际生产可以两种表都建流水表用于追溯和重建镜像表用于日常查询。两张表共用同一条 CDC 链路只是写入逻辑不同。5. 环境准备与前置条件下面给出一条通用的准备清单具体路径和版本按你的实际环境调整。5.1 MySQL 侧确认 MySQL 已开启 binlog并检查参数SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;建议配置server_id 223344 log_bin mysql-bin binlog_format ROW binlog_row_image FULL expire_logs_days 7 gtid_mode ON enforce_gtid_consistency ON其中 binlog_row_image 设置为 FULL可以保证更新事件里同时有完整的 before 和 after 数据。expire_logs_days 要根据链路消费速度设置太短会导致下游追不上时 binlog 已被清理太长会占磁盘。创建 CDC 专用账号并授权CREATE USER cdc_user% IDENTIFIED BY your_password; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES;需要说明REPLICATION SLAVE 权限是 Debezium、Canal 等工具读取 binlog 的常用要求具体看所选工具文档如果只需要读取业务表数据做全量初始化保留 SELECT 权限即可。5.2 BigQuery 侧提前创建数据集和表。BigQuery 是列式存储表结构最好按目标业务字段显式定义不要每次写入时自动推断避免类型漂移。命令行创建示例bq mk --project_idyour_project your_dataset bq query --use_legacy_sqlfalse \ CREATE TABLE IF NOT EXISTS your_dataset.orders_sink ( order_id STRING, status STRING, amount NUMERIC, updated_at TIMESTAMP, op_type STRING, binlog_ts TIMESTAMP )如果使用 Storage Write API需要保证服务账号有bigquery.tables.updateData权限。5.3 工具部署环境不管选哪种 CDC 组件都建议准备 Docker 或独立服务器。下面几条是通用检查项网络MySQL 允许 CDC 组件所在主机访问 3306 端口。磁盘binlog 会占磁盘CDC 组件所在机器也要预留日志和本地缓冲空间。JavaCanal、Debezium、Flink 都依赖 JDK建议按所选组件要求安装指定 JDK 版本。时间同步源库、CDC 组件、BigQuery 之间时区不一致会出现时间偏移建议统一使用 UTC 并在应用层转换。6. CDC 链路搭建与运行验证本节给两条可操作的路线一条是 Flink CDC 的 SQL/YAML 模式适合从零快速跑通另一条是 Debezium Kafka 的模式适合已有消息队列的场景。6.1 快速跑通Flink CDC 同步到 BigQuery前提是本机已安装 Flink 或能运行 Flink SQL 客户端且已经下载对应的 MySQL CDC connector 和 BigQuery connector。以下 YAML 配置是 Flink CDC 常见的本地作业定义实际连接信息需要替换source: type: mysql hostname: 127.0.0.1 port: 3306 username: cdc_user password: your_password tables: app_db.orders server-id: 5400-5404 server-time-zone: Asia/Shanghai scan.incremental.snapshot.enabled: true scan.startup.mode: initial sink: type: bigquery projectId: your_project dataset: your_dataset table: orders_sink writeMethod: STORAGE_WRITE_API启动 Flink CDC 作业后在 MySQL 里执行几条 INSERT、UPDATE、DELETE再到 BigQuery 查询目标表应该能看到事件已经写入。判断成功的标准新插入的行能在秒级出现在 BigQuery。UPDATE 后目标表对应主键的值发生变化。DELETE 事件在流水表中出现或者镜像表里对应行被标记删除。启动时如果遇到“找不到连接器”的报错先确认 connector jar 是否放在 Flink 的 lib 目录下。6.2 生产常用Debezium KafkaDebezium 的 MySQL Connector 以 Kafka Connect 的方式运行。部署时先启动 Kafka 和 Kafka Connect再提交连接器配置。连接器配置示例{ name: mysql-orders-connector, config: { connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: 127.0.0.1, database.port: 3306, database.user: cdc_user, database.password: your_password, database.server.id: 223344, database.server.name: mysql-server, database.include.list: app_db, table.include.list: app_db.orders, database.history.kafka.bootstrap.servers: kafka:9092, database.history.kafka.topic: schema-changes.app_db, include.schema.changes: true, tombstones.on.delete: false, key.converter: org.apache.kafka.connect.json.JsonConverter, value.converter: org.apache.kafka.connect.json.JsonConverter } }Debezium 输出的 Kafka 消息中value 里主要字段包括 before、after、source、op。op 的取值一般是 c新增、u更新、d删除、r初始快照读取。下游服务消费 Kafka 消息后再写入 BigQuery。一种简单的 Python 消费端模板import json from google.cloud import bigquery from kafka import KafkaConsumer client bigquery.Client() table_id your_project.your_dataset.orders_sink consumer KafkaConsumer( mysql-server.app_db.orders, bootstrap_servers127.0.0.1:9092, value_deserializerlambda m: json.loads(m.decode(utf-8)) ) for message in consumer: event message.value after event.get(after) or {} op event.get(op) row { order_id: after.get(order_id), status: after.get(status), amount: after.get(amount), op_type: op, } errors client.insert_rows_json(table_id, [row]) if errors: print(errors)这段代码只是演示批量消费写入的骨架。生产环境要看情况加上幂等键去重、错误重试、批量攒批等逻辑避免单条插入导致吞吐过低。6.3 存量数据初始化CDC 增量同步只能接住“开启后发生的新变化”存量数据需要先导入 BigQuery。Flink CDC 的scan.startup.mode: initial会在启动时先做一次全量快照再继续消费 binlog这就是“存量 增量”的最简单实现。如果自研链路可以在启动 CDC 前先导出 MySQL 当前数据再启动增量最后用这个时间点作为拼接边界。7. 接口 API 与批量任务设计7.1 CDC 组件本身的接口Canal 提供基于 TCP 和 RocketMQ 的客户端接口也提供 HTTP 管理端口用于查看消费位点Kafka Connect 提供 REST API可以提交连接器、查看任务状态。这些接口方便运维不需要登录服务器看日志。常用操作示例# 查看 Kafka Connect 上的连接器状态 curl http://127.0.0.1:8083/connectors/mysql-orders-connector/status # 重启连接器 curl -X POST http://127.0.0.1:8083/connectors/mysql-orders-connector/restart如果你希望把变更事件直接通过 HTTP 推给下游系统可以做一个轻量订阅服务收到 Kafka 消息后调用目标系统 API。这里不涉及具体平台限制核心是给事件流加一个可控的出口。7.2 BigQuery 写入方式与批量任务BigQuery 写入有两条主要路径Storage Write API适合高吞吐流式写入支持 exactly-once 语义是 CDC 推荐的写入方式。批量加载Load Job适合每天把文件导入分区表CDC 里一般用于存量初始化或周期性对账。事务性任务设计时要关注“批流关系”。最佳实践是把 CDC 流式写入作为主链路把每日定时任务作为“校验和补偿”。每日任务扫描 BigQuery 中的镜像表和 MySQL 做一次主键比对发现漏掉的变更再补发。这么设计的好处是流式写入负责低延迟批量任务负责兜底不会因为某一条消息丢失造成永久差异。7.3 批量回放binlog 天然支持回放。因为事件里记录了变更前后的行数据只要把这些事件按 gtid 或 binlog 位点排序再重新写入 BigQuery就能修复下游表。做回放前注意事件要按事务边界分组不要拆散同一个事务的多个事件。回放时目标表需要有主键或去重逻辑。先在一个测试数据集里回放确认数据量一致后再切换正式环境。8. 资源占用与性能观察8.1 MySQL binlog 的额外开销开启 binlog 本身会有写入开销ROW 格式下一条 UPDATE 如果涉及多行会产生多组 before/after 数据binlog 文件增长速度会明显加快。观察指标binlog 文件总大小。单位时间 binlog 增量。磁盘剩余空间。如果 binlog 增长太快、消费端跟不上排查顺序通常是先看消费端是否积压再看是否有大事务一次性修改大量数据最后确认 binlog 清理策略是否合理。8.2 CDC 链路的内存与磁盘占用Canal 和 Debezium 在启动全量快照时会比较吃内存和临时磁盘因为要把大规模查询结果做成事件。Flink CDC 做全量快照时也会占用 state。建议给这些组件至少预留 4G 以上内存特别是在同步 10 万行以上的大表时。8.3 BigQuery 写入配额BigQuery 对写入请求数量和单表分区写入有配额限制。CDC 链路如果事件量很大建议在写入端做攒批而不是一条一条调用。同时可以把大表按日期分区避免写入压力都集中到同一个分区。8.4 性能观察清单观察对象关注指标常见问题MySQLbinlog 增长速率、慢查询binlog 被清理导致消费端重连失败CDC 组件消费延迟、堆内存、线程数大事务造成解析暂停Kafkatopic 积压量、消费速率下游写入慢造成消息积压BigQuery写入行数、失败请求数RPC 超时或限流9. 常见问题与排查方法问题现象可能原因排查方式解决方案CDC 启动后立即报“binlog not enabled”MySQL 未开启 binlog执行 SHOW VARIABLES LIKE log_bin开启 binlog 并重启 MySQL停止一段时间后重启无法找到 binlog 位点binlog 被 expire_logs_days 清理查看 MySQL 当前 binlog 文件列表缩短消费停顿时间或先做全量初始化再增量Flink CDC 运行时提示 server-id 冲突多个进程使用相同 server_id查看连接器配置给每个作业配置不同 server-id 范围BigQuery 中时间比源库早/晚 8 小时时区处理不一致检查 JDBC 连接时区、Flink 时区参数统一使用 UTC 或明确设置 server-time-zoneMySQL 大事务执行后事件延迟暴增单个事务修改大量行CDC 需逐个事件解析查看 binlog 中该事务事件数量拆分大事务或提高 CDC 并发删除事件在 BigQuery 里看不到表模型是覆盖式删除未映射检查 op 为 d 的事件是否被过滤在目标表增加 op_type 字段或使用标记删除写入 BigQuery 出现 quota exceeded单分区写入次数过高查看 BigQuery 监控按主键分布分区写入端攒批这一节是实操中最高频的几个坑。遇到问题第一步永远是看日志不要直接改配置。Canal、Flink CDC、Debezium 的日志里都会输出解析到哪个 binlog 位点这个信息对整个链路排错最有价值。10. 最佳实践与合规建议10.1 工程化建议第一次跑通时先用小表不要上来就同步几百 GB 的大表。把 binlog 消费位点持久化。不管是 Kanel、Canal 还是 Flink消费位点丢失会导致重复或缺失。给 CDC 组件单独建账号权限限制在业务库的 SELECT 和复制权限不要用 root。对目标表做分区设计。BigQuery 按时间分区可以降低查询成本和单分区写入压力。对事件消息做幂等键。例如用source_table primary_key binlog_position生成唯一 ID防止重复消息造成目标表数据重复。保留一段时间原始事件。不要只把字段拆开写入目标表建议同时保留完整 JSON 事件到另一张表便于排查数据质量问题。批量补数任务和 CDC 主链路并行时要设计好重叠窗口避免同一行被两套逻辑写入。10.2 数据合规与安全边界CDC 会把 MySQL 里几乎所有数据变更都搬到 BigQuery链路一旦打通数据生命周期就变长了。使用时必须注意明确同步范围只同步业务需要的表不要默认“把所有库同步过去”。涉及用户身份证、手机号、地址等敏感字段时在链路中做脱敏或掩码处理。授权边界要落到具体账号、数据集、表级别。对增量事件和审计日志设置合理的保留周期。如果同步的是第三方授权数据或公版素材先确认授权范围允许复制和存储。在生产环境上线前对同步链路做一次最小权限测试和隐私影响评估。CDC 不是把数据“多复制一份”这么简单它实际上是把数据库内部操作永久记录下来。因此合规策略必须和同步链路同时设计而不是等功能上线后再补。11. 小结与下一步MySQL CDC 到 BigQuery 最核心的转变是把“定时全量拉取”换成“binlog 事件流持续消费”。它解决的不只是延迟问题更重要的是把删除、更新、过程变化和 DDL 行为完整地保留下来。对需要实时报表、数据审计、跨系统一致性的团队来说这条链路值得尽早搭建。如果今天只做一件事先把 MySQL 的 binlog 配置检查一遍再用 Flink CDC 或 Canal 跑通一张小表的增量同步观察 BigQuery 里 insert、update、delete 三类事件是否都正确反映。跑通后再考虑大表、多表、分布式扩展和批量补偿。最容易踩的坑还是 binlog 清理时机、server-id 冲突、时区不统一和 BigQuery 写入配额建议收藏备用等真上线时逐项核对。
返回列表