拓冰建站拓冰建站
首页 / 资讯中心 / 正文

MySQL CDC 到 BigQuery:binlog 与 Debezium+Kafka 实时同步实战

做数据同步的人迟早会遇到同一个问题业务方过来问“为什么数仓里的订单数和线上差那么多”而你只能回答“下次同步还没跑”“延迟半小时”。这种对话出现得越来越频繁往往意味着定时同步方案已经撑不住业务体量了。本文围绕 MySQL CDCChange Data Capture变更数据捕获到 BigQuery 的同步链路展开重点对比周期性同步的盲区并通过 binlog 讲清楚 CDC 的底层原理最后给出一套可复现的 Debezium Kafka 实验环境。不管你是正在维护数据管道的数据工程师还是想把业务库数据搬到数仓的后端开发者都能从这篇文章里找到可以直接落地的思路。文章中不会只讲概念也不会只贴代码每一步都会解释“为什么这么做”以及生产环境里容易踩的坑。1. 周期性同步到底丢掉了什么1.1 定时同步的两种常见实现在引入 CDC 之前大多数团队的同步方案是“定时跑脚本”。思路通常有两种。第一种是时间戳增量同步。假设源表里有一个update_time字段脚本每隔几分钟或几小时执行一次读取大于上次同步时间点的数据写入目标库。# 伪代码每 5 分钟执行一次 def sync_orders(): last_time get_last_sync_time() with source_db.cursor() as cur: cur.execute( SELECT * FROM orders WHERE update_time %s, last_time ) rows cur.fetchall() for row in rows: upsert_to_bigquery(row) set_last_sync_time(now())第二种是全量比对同步。把源表和目标表的数据全部拉下来逐条比较差异再统一更新。这种方案在小表阶段很好用但一旦数据量上去了每次跑全量都是对源库和目标库的双重压力。# 伪代码全量比对 def full_sync_orders(): source_rows load_all_source_orders() target_rows load_all_target_orders() diff compare_two_sets(source_rows, target_rows) apply_diff(diff)这两种方式看起来都能工作但它们在数据一致性上存在天然盲区。1.2 三个最容易翻车的场景先给出一张典型的订单表结构后面所有讨论都围绕这张表展开CREATE TABLE orders ( id BIGINT PRIMARY KEY AUTO_INCREMENT, order_no VARCHAR(32) NOT NULL, user_id BIGINT NOT NULL, amount DECIMAL(10, 2) NOT NULL, status VARCHAR(20) NOT NULL, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, KEY idx_status (status) );场景一物理删除无法感知。假如业务方误插入了一批测试订单或者在风控流程里删除了异常订单数据库执行的是DELETE FROM orders WHERE id 1001。时间戳增量同步完全感知不到这次删除目标库 BigQuery 里那条脏数据会一直存在。有人会说“我们业务不做物理删除只做软删除”。但实际情况是总有一些表由于历史原因没有软删除字段或者 DBA 手动清理数据这类例外恰恰是数据质量事故的高发区。场景二业务字段更新不规范。时间戳增量同步有一个强依赖业务代码在 UPDATE 时必须正确维护update_time字段。但现实是很多 UPDATE 语句并不会主动带上这个字段。-- 没有更新 update_time甚至直接改主键 UPDATE orders SET status PAID WHERE order_no A001;如果有人执行了上面的语句或者直接用工具连库修改了几条老数据增量同步下一次扫描时根本扫不到这些变化。数据不一致就这样悄悄发生了。场景三源库压力不可控。全量比对方案每次都要全表扫描索引再优化也扛不住每天几十次扫描大促、月末对账期间这种同步脚本还会和业务抢数据库资源导致慢查询变多甚至拖慢线上接口。对数据量大的表来说周期同步的时间窗会越来越长最终从“每 5 分钟一次”退化到“每小时一次”甚至“明天再说”。2. CDC 是什么binlog 为什么是关键2.1 CDC 的核心思想CDC 全称 Change Data Capture中文叫“变更数据捕获”。它的核心思想是不主动去查“哪些数据变了”而是让数据库主动告诉下游“哪一行数据发生了变化”。打个比方周期性同步像保安每隔半小时去仓库清点一遍货数完才能知道少了什么而 CDC 相当于在仓库门口装了一个感应器有人进出就会触发记录。两者获取信息的及时性和准确度完全不同。实现 CDC 的技术方案有很多种MySQL 生态里最常见的做法是解析 binlog。2.2 binlog 是 MySQL 的“操作日志”binlog 是 MySQL Server 层生成的二进制日志记录了所有导致数据变化的操作。主从复制就是基于 binlog 实现的主库把 binlog 发给从库从库重放这些日志从而保持数据一致。binlog 有三种格式这一点在配置 CDC 时非常关键STATEMENT语句格式记录的是 SQL 原文。比如UPDATE orders SET statusPAID数据量小但主从不一致风险高。ROW行格式记录的是每行数据“变化前”和“变化后”的具体内容。信息最完整最适合 CDC。MIXED混合格式默认使用 STATEMENT遇到不确定操作时自动切换为 ROW。对于 CDC 工具来说必须把 binlog 设置为 ROW 格式并且把binlog_row_image设置为FULL才能拿到每一行数据变化前后的完整镜像。2.3 CDC 工具如何“伪装成从库”MySQL 主从复制的模型是从库连接主库请求指定 binlog 位置之后的日志然后重放。CDC 工具的工作原理与此类似工具把自己伪装成一个“从库”向 MySQL 注册一个唯一的server-id请求从某个 binlog 位置开始同步。MySQL 把 binlog 事件源源不断地发给它工具解析这些事件转换成结构化的消息再发给下游。正因为 binlog 是 MySQL 原生日志CDC 方案对源库业务代码几乎零侵入。不需要业务方改表结构不需要统一维护update_time也不需要每张表都做软删除改造。3. 主流 MySQL CDC 工具怎么选3.1 Debezium Kafka ConnectDebezium 是目前社区最活跃的 CDC 框架之一。它基于 Kafka Connect 运行把 MySQL binlog 事件转换成标准化的 JSON 或 Avro 消息写入 Kafka。这种方案的优势很明显Kafka 作为中间缓冲层可以解耦源库和目标库消息可以重复消费方便多个下游使用同一份变更数据同时 Kafka 本身具备消息堆积能力能应对下游短时故障。3.2 Flink CDCFlink CDC 是 Flink 生态的连接器底层同样通过 Debezium 解析 binlog但把能力封装成了 Flink 的 Source。它最大的优势是能和 Flink SQL、Flink 流处理任务无缝集成适合需要做实时计算、实时关联的场景。3.3 CanalCanal 是阿里巴巴开源的项目专门针对 MySQL binlog 解析。它的诞生最初是为了解决业务缓存同步问题因此很多 Java 技术栈的团队用得比较多。Canal 可以将 binlog 发送到 MQ如 RocketMQ、Kafka或者直接推送给下游应用。3.4 选型建议选型没有唯一答案主要看现有技术栈如果团队已经有 Kafka并且希望下游有多方消费优先选Debezium Kafka Connect。如果团队重度使用 Flink需要边同步边计算优先选Flink CDC。如果团队是 Java 技术栈数据链路简单Canal也完全足够。本文后续的实战部分以 Debezium 为主因为它最接近“binlog 订阅与分发”的本质也方便把每一步拆开看清楚。4. 环境准备开启 MySQL binlog4.1 确认当前 binlog 状态在开始之前先确认 MySQL 是否已经开启了 binlog。使用 MySQL 客户端登录后执行SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;如果log_bin的值为OFF需要修改配置并重启 MySQL。云数据库如阿里云 RDS、腾讯云 CDB通常会默认开启 binlog但格式不一定满足 CDC 要求需要根据控制台文档确认。4.2 开启 binlog 并配置 ROW 格式编辑 MySQL 配置文件/etc/my.cnf或自定义的 MySQL Docker 配置在[mysqld]段添加以下内容[mysqld] server_id 223344 log_bin mysql-bin binlog_format ROW binlog_row_image FULL如果使用 MySQL 5.7可以设置 binlog 保留时间expire_logs_days 14如果使用 MySQL 8.0推荐使用秒数配置binlog_expire_logs_seconds 1209600这里要注意server_id必须全局唯一。如果多台机器连接同一个 MySQL 实例或者多个 CDC 任务连接到同一个实例它们的server_id不能重复否则会导致连接断开。修改配置后重启 MySQLsystemctl restart mysqld4.3 创建 CDC 专用账号为了安全不要用 root 账号直接采集 binlog。建议创建专用账号只授予最小权限。CREATE USER cdc_user% IDENTIFIED BY cdc_pass; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES;这里每个权限都有自己的用途SELECTDebezium 做初始快照时需要读取表数据。REPLICATION SLAVE允许伪装成从库请求 binlog。REPLICATION CLIENT允许查询主库状态例如SHOW MASTER STATUS。RELOAD某些快照场景需要执行FLUSH TABLES WITH READ LOCK保证一致性。4.4 验证 binlog 是否产生重新登录 MySQL执行SHOW MASTER STATUS;正常情况下会看到类似下面的输出说明 binlog 已开启------------------------------------------------------------------------------- | File | Position | Binlog_Do_DB | Binlog_Ignore_DB | Executed_Gtid_Set | ------------------------------------------------------------------------------- | mysql-bin.000003 | 156 | | | | -------------------------------------------------------------------------------执行几条 DML 语句后再用SHOW BINLOG EVENTS IN mysql-bin.000003 LIMIT 5;查看就能看到 binlog 事件了。5. 实战Debezium 采集 MySQL 变更到 Kafka5.1 用 Docker Compose 搭一套最小环境为了方便复现这里直接用 Docker Compose 启动 MySQL、Zookeeper、Kafka、Debezium Connect 四个组件。创建docker-compose.ymlversion: 3.8 services: mysql: image: mysql:8.0 container_name: mysql environment: MYSQL_ROOT_PASSWORD: root MYSQL_DATABASE: shop command: - --server-id223344 - --log-binmysql-bin - --binlog-formatROW - --binlog-row-imageFULL - --binlog-expire-logs-seconds1209600 ports: - 3306:3306 volumes: - mysql-data:/var/lib/mysql zookeeper: image: confluentinc/cp-zookeeper:latest container_name: zookeeper environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: confluentinc/cp-kafka:latest container_name: kafka depends_on: - zookeeper environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 ports: - 9092:9092 connect: image: debezium/connect:latest container_name: connect depends_on: - kafka environment: BOOTSTRAP_SERVERS: kafka:9092 GROUP_ID: 1 CONFIG_STORAGE_TOPIC: connect-configs OFFSET_STORAGE_TOPIC: connect-offsets STATUS_STORAGE_TOPIC: connect-status KEY_CONVERTER: org.apache.kafka.connect.storage.StringConverter VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter VALUE_CONVERTER_SCHEMAS_ENABLE: false ports: - 8083:8083 volumes: mysql-data:启动服务docker-compose up -d生产环境强烈建议锁定镜像 tag而不是使用latest。这里使用latest只是为了演示方便。5.2 准备测试表和数据进入 MySQL 容器docker exec -it mysql mysql -uroot -proot shop创建订单表CREATE TABLE orders ( id BIGINT PRIMARY KEY AUTO_INCREMENT, order_no VARCHAR(32) NOT NULL, user_id BIGINT NOT NULL, amount DECIMAL(10, 2) NOT NULL, status VARCHAR(20) NOT NULL, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL ); INSERT INTO orders (order_no, user_id, amount, status, create_time, update_time) VALUES (A001, 1001, 99.50, CREATED, NOW(), NOW());5.3 注册 Debezium Connector通过 Kafka Connect 的 REST API 注册 MySQL Connectorcurl -i -X POST \ -H Accept:application/json \ -H Content-Type:application/json \ http://localhost:8083/connectors/ \ -d { name: mysql-orders-connector, config: { connector.class: io.debezium.connector.mysql.MySqlConnector, database.hostname: mysql, database.port: 3306, database.user: cdc_user, database.password: cdc_pass, database.server.id: 223345, database.server.name: mysql-demo, database.include.list: shop, table.include.list: shop.orders, schema.history.internal.kafka.topic: schema-changes.mysql, snapshot.mode: initial } }这里有几个配置需要重点说明database.server.idDebezium 连接 MySQL 时使用的 server-id必须与源库的server_id不同并且与其他连接器保持唯一。database.server.name决定 Kafka topic 的前缀。默认情况下topic 名称为{server.name}.{database}.{table}。schema.history.internal.kafka.topicDebezium 会把历史的 schema 信息存储在 Kafka 里这样后续解析老 binlog 时还能知道当时的表结构。如果使用 Debezium 1.x配置项名称是database.history.kafka.topic。等待几十秒后检查连接器状态curl http://localhost:8083/connectors/mysql-orders-connector/status看到state: RUNNING说明连接器已正常运行。5.4 观察 insert / update / delete 消息Debezium 启动时会先执行一次一致性快照已有数据会以op: rread的消息写入 Kafka。随后源库的增删改操作会实时产生消息。先在 MySQL 中插入一条新数据INSERT INTO orders (order_no, user_id, amount, status, create_time, update_time) VALUES (A002, 1002, 29.90, CREATED, NOW(), NOW());启动 Kafka 自带的消费者查看消息docker exec -it kafka kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic mysql-demo.shop.orders \ --from-beginning \ --max-messages 5插入操作对应的消息大致如下这里简化了 schema 部分{ before: null, after: { id: 2, order_no: A002, user_id: 1002, amount: 29.90, status: CREATED, create_time: 1690000000000, update_time: 1690000000000 }, source: { db: shop, table: orders, ts_ms: 1690000001000 }, op: c }再执行一次更新UPDATE orders SET status PAID, update_time NOW() WHERE order_no A002;更新操作对应的消息op变为u同时before和after都会存在{ before: { id: 2, order_no: A002, user_id: 1002, amount: 29.90, status: CREATED, create_time: 1690000000000, update_time: 1690000000000 }, after: { id: 2, order_no: A002, user_id: 1002, amount: 29.90, status: PAID, create_time: 1690000000000, update_time: 1690000001000 }, source: { db: shop, table: orders, ts_ms: 1690000002000 }, op: u }关键是删除操作。执行DELETE FROM orders WHERE order_no A002;即使目标库没有做任何特殊处理删除消息也会被 binlog 完整记录{ before: { id: 2, order_no: A002, user_id: 1002, amount: 29.90, status: PAID, create_time: 1690000000000, update_time: 1690000001000 }, after: null, source: { db: shop, table: orders, ts_ms: 1690000003000 }, op: d }这就是 binlog 方案与周期同步最大的区别删除操作不再丢失更新也不依赖业务方是否维护update_time字段一切以数据库真实发生的变更为准。5.5 消息里每个字段的含义Debezium 的消息体结构是固定的理解这些字段对下游处理非常有帮助before数据变更前的镜像。插入操作中为null更新和删除操作中可能存在。after数据变更后的镜像。删除操作中为null。op操作类型。c表示插入u表示更新d表示删除r表示初始快照阶段读取。source.ts_ms变更事件发生的时间戳毫秒通常对应 binlog 中记录的时间。source.db/source.table变更事件所在的数据库和表。下游拿到这些消息后可以在目标侧执行 upsert、逻辑删除、物理删除等不同的数据策略灵活度比周期同步高很多。6. 从 Kafka 接入 BigQuery 的三条路线把 MySQL 的变更数据送到 Kafka 只是第一步接下来还需要将数据落地到 BigQuery。取决于团队的基础设施有三条比较常用的路线。6.1 路线一Kafka Connect BigQuery Sink社区中有开源的 Kafka Connect BigQuery Sink 连接器常见的实现来自 GitHub 上的kafka-connect-bigquery项目。它的工作方式是在 Kafka Connect 集群中再部署一个 Sink 连接器消费mysql-demo.shop.orders主题把 JSON 消息写入 BigQuery。注册 BigQuery Sink 的配置大致如下不同版本参数会有差异请以项目 README 为准{ name: bigquery-sink, config: { connector.class: com.wepay.kafka.connect.bigquery.BigQuerySinkConnector, topics: mysql-demo.shop.orders, projectId: your-gcp-project, datasets: ordersshop_ods, keyfile: /etc/bigquery/keyfile.json, autoCreateTables: true, sanitizeTopics: true } }这种方案适合已经使用 Kafka Connect 的团队复用现有的连接器机制运维成本相对低。需要注意该连接器默认会使用 BigQuery 的流式插入接口对数据量很大的场景需要考虑配额和费用。6.2 路线二Kafka → Pub/Sub → Dataflow → BigQuery这是 GCP 原生生态里很常见的管道设计Kafka 中的数据通过 Pub/Sub 连接器转发到 Pub/Sub 主题然后由 Dataflow 作业消费 Pub/Sub 消息经过必要的清洗和转换后写入 BigQuery。这条链路看起来多了几个组件但优势在于一旦数据进入 Pub/Sub后续就有了 GCP 托管的流处理能力。Dataflow 自带的窗口、去重、流式写入 BigQuery 语义都比较成熟适合对数据延迟和准确性要求较高的生产链路。6.3 路线三Flink CDC 直连 BigQuery Connector如果不想引入 Kafka 中间层也可以直接用 Flink CDC 从 MySQL 采集数据再通过 Flink 的 BigQuery Sink 写入。Flink CDC 的源表定义非常简洁下面是一个可运行的 Flink SQL 示例CREATE TABLE orders_cdc ( id BIGINT, order_no STRING, user_id BIGINT, amount DECIMAL(10, 2), status STRING, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username cdc_user, password cdc_pass, database-name shop, table-name orders, server-id 5400 );之后通过INSERT INTO语句将数据写入 BigQuery 目标表。Flink CDC 底层同样依赖 Debezium 解析 binlog因此它具备 Debezium 的实时性和变更捕获能力同时把数据处理逻辑放进同一条流中。6.4 BigQuery 侧的前置准备无论使用哪条路线访问 BigQuery 前都需要完成以下准备工作在 GCP 中创建项目并启用 BigQuery API。创建服务账号授予 BigQuery Data Editor 或更细粒度的权限。下载服务账号的 JSON Key部署到 Kafka Connect 或 Dataflow 所在节点。在 BigQuery 中提前创建数据集和表。BigQuery 侧的目标表可以直接参考 MySQLorders表的结构来建。如果对字段类型不敏感也可以让连接器自动建表但生产环境更推荐由 DBA 和数据工程师统一设计明确分区字段和 clustering 字段。7. 周期性同步 vs CDC一个对比表为了更直观地说明两种方案差异这里做一个汇总对比维度周期性同步binlog / CDC 方案时效性分钟级、小时级或 T1秒级几乎实时删除捕获很难依赖全量比对天然捕获binlog 记录删除事件依赖业务字段强依赖update_time等字段不依赖读取数据库原生日志源库压力每次扫描消耗大量资源解析日志压力很小数据一致性存在明显不一致窗口变更即时传递一致性窗口极短历史数据追溯难以追溯binlog 保留期内可回溯实现复杂度简单脚本即可需要维护 CDC 工具和消息通道schema 变更需要手动感知连接器会记录 schema 历史表格里没有绝对的“谁更好”因为周期性同步在小数据量、非核心场景下仍有成本优势。但当数据量和实时性要求上来之后binlog 方案的优势会越来越明显。8. 常见问题与排查思路8.1 常见错误对照表以下是 CDC 链路搭建过程中最常遇到的几类问题问题现象常见原因解决思路连接器启动报错Access deniedCDC 账号缺少REPLICATION SLAVE权限按上文重新授予权限连接器反复重连server-id与已有连接重复为每个连接器分配全局唯一server-id消息里只有after没有beforebinlog_row_image不是FULL修改配置并重启 MySQL断连后无法接续同步binlog 已被清理调整 binlog 保留时间测试重连恢复下游数据延迟越来越大Kafka 消费能力不足或目标库写入慢增加 Sink 任务并发检查下游负载UPDATE 消息丢失源库存在大事务binlog 事件积压拆分大事务优化 CDC 任务资源8.2 三个高频场景排查场景一快照完成但增量不生效。检查连接器的snapshot.mode。默认initial表示先做快照再读增量。如果表数据量很大快照阶段会比较长这期间的新增变更不会丢失Debezium 会在快照完成后继续从记录的位置读取 binlog。如果设置了schema_only则不会复制已有数据只有后续增量。场景二topic 里出现大量d操作目标库不想物理删除。删除消息在目标侧可以有两种处理方式直接删除目标行或者做逻辑删除例如在目标表增加is_deleted字段并标记。具体策略需要和目标表的下游消费方确认。场景三binlog 消费断点过期。假设 CDC 任务停了三天而 MySQL 的 binlog 只保留两天重新启动后会出现“binlog 文件找不到”的报错。这种情况只能重新做一次全量快照。因此生产环境的 binlog 保留时间必须覆盖“故障恢复时间 下游排查时间”建议配置 7 天以上并根据数据量评估磁盘成本。9. 落地建议与下一步把周期性同步换成 CDC不是技术炫技而是当数据体量和实时性要求上来之后早晚要迈过去的一步。建议先从一张非核心订单表开始实验按照本文第四、五章的步骤搭建最小链路验证 insert、update、delete 三类消息都能正常到达下游再逐步推广到核心表。具体落地时有几点可以优先确认先保证 binlog 保留时间合规。这是最容易在接入后期踩坑的问题建议上线前就和 DBA 确认清楚。账号权限遵循最小化原则。CDC 账号只给 SELECT 和复制相关权限不要使用 root。监控要跟上。重点监控 Kafka 消费延迟、BigQuery 写入失败率、MySQL 连接数变化。定期做故障演练。停掉 Kafka 一段时间再恢复观察 CDC 任务能否从 binlog 断点继续消费这比任何文档都更有说服力。下一步可以继续深入了解 Debezium 的 schema 变更处理机制、Flink CDC 的 checkpoint 语义以及 BigQuery 流式写入的成本控制。数据同步链路只要跑起来就会出现新的问题把每个问题记录成排查笔记是工程师成长最快的方式。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门