CDC 与实时数据管道:Canal、Debezium、Flink CDC 的架构与实战

业务库的数据是实时的,而下游的数据仓库、缓存、搜索引擎往往是滞后的——传统 ETL 定时全量同步的延迟以"小时"计,且全量扫描对源库冲击大。CDC(Change Data Capture,变更数据捕获)改变了这一切:它从数据库的日志(MySQL binlog、PostgreSQL …

业务库的数据是实时的,而下游的数据仓库、缓存、搜索引擎往往是滞后的——传统 ETL 定时全量同步的延迟以"小时"计,且全量扫描对源库冲击大。CDC(Change Data Capture,变更数据捕获)改变了这一切:它从数据库的日志(MySQL binlog、PostgreSQL WAL、MongoDB oplog)里"偷听"每一次变更,以毫秒级延迟把增量数据流式推给下游。本指南深入 CDC 的底层原理(binlog/WAL 解析、位点管理、Exactly-Once)、主流工具(Canal、Debezium、Flink CDC)的架构与选型、与 Kafka/Flink 的组合链路、Schema 演进处理,以及生产级 CDC 管道的设计与避坑。

一、CDC 是什么:从"拉取"到"订阅"

1.1 传统同步 vs CDC

传统 ETL(拉取):
  定时全量 SELECT → 传输 → 写入
  · 延迟:小时级
  · 成本:全表扫描,冲击源库
  · 捕获不到"删了哪些行/改了什么"(只有当前状态)

CDC(订阅):
  解析数据库日志 → 逐条变更事件 → 实时推送
  · 延迟:毫秒级
  · 成本:只读日志,几乎零影响源库
  · 捕获完整变更(INSERT/UPDATE/DELETE + before/after)

ℹ️ 核心洞察:CDC 把数据库变成"事件源"——业务里的每一次写操作,都天然成为下游系统可以订阅的数据事件流。这是实时数仓、缓存同步、搜索索引、审计追踪的统一底座。

1.2 CDC 的典型场景

· 实时数仓:业务库 → CDC → Kafka → Flink → DWS/DW
· 缓存刷新:订单变更 → 实时刷新 Redis 缓存
· 搜索索引:商品变更 → 实时同步到 ES
· 微服务数据共享:订单库变更 → 同步给下游服务
· 审计与合规:记录所有变更(before/after)
· 数据迁移/双写过渡

二、CDC 的底层原理:binlog 与 WAL

2.1 MySQL binlog

binlog 三种格式:
  STATEMENT:记录 SQL(节省空间,但非确定性函数有风险)
  ROW      :记录每行变更前后值(推荐,信息最全)
  MIXED    :混合,默认 statement 必要时转 row

CDC 必须用 ROW 格式:
  ROW 事件包含 before_image / after_image
  → 才拿得到完整的字段变化

binlog 三种日志记录方式:
  · 简单轮询读取日志文件
  · 基于 binlog 位点(position)断点续传(Canal 传统方式)
  · GTID(全局事务 ID):更可靠的断点恢复(推荐)
-- 确认 binlog 配置
SHOW VARIABLES LIKE 'binlog_format';        -- 需为 ROW
SHOW VARIABLES LIKE 'gtid_mode';            -- 建议 ON
SHOW VARIABLES LIKE 'log_bin';              -- 需 ON

2.2 PostgreSQL WAL / 逻辑复制

PG 的 CDC 走"逻辑复制":
  · 物理复制:WAL 原始字节(用于流复制,对库不可读)
  · 逻辑复制:WAL 解码为 SQL 语义的变更流(CDC 用)
  · 通过 pgoutput / wal2json 插件输出 JSON 变更

PG CDC 特点:
  · 发布/订阅模式(PUBLICATION/SUBSCRIPTION)
  · 天然支持 schema 变化与过滤
  · Debezium 等工具消费逻辑复制输出

2.3 位点(Offset)管理

CDC 必须记录"读到哪了"(offset/position/GTID):
  · 进程重启 → 从上次位点继续,不丢不重
  · 存到 Kafka offset / state store / DB

Exactly-Once 挑战:
  · 源日志消费(offset)与下游写入(Kafka/DB)跨系统
  → 需"幂等 + 去重"或依赖 Flink checkpoint 的端到端一致

三、主流工具架构

3.1 Canal(阿里,MySQL → 消息队列)

架构:
  Canal Instance 伪装成 MySQL 从库
    → 拉取 binlog
    → 解析为结构化变更事件(CanalEntry)
    → 推给下游(Kafka / RocketMQ / DB)

适用:
  · 纯 MySQL → Kafka/MQ 的轻量管道
  · 已有 Kafka 生态的团队
  · 对 Debezium 生态陌生、想更贴近阿里实践
# canal.properties / instance.properties
canal.instance.master.address=mysql-master:3306
canal.instance.dbUsername=canal
canal.instance.dbPassword=***
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=orders_db\\..*
# 位点持久化
canal.instance.gtidon=false
canal.instance.master.position=

3.2 Debezium(Red Hat,跨库 CDC 框架)

架构:
  Debezium Connector(Kafka Connect 插件)
    → 连接源库(MySQL/PG/MongoDB/Oracle/SQL Server...)
    → 解析日志 → 输出到 Kafka topic(每个表一个 topic)
    → Kafka Connect 管理位点(offset storage)

特性:
  · 原生 Kafka 生态(Avro/JSON、schema registry 集成)
  · 内置字段 flatten、过滤、列裁剪
  · 支持 schema 变更事件(DDL)
  · Exactly-Once 依赖 Kafka 事务 / Flink
// Debezium MySQL connector 配置(Kafka Connect REST API)
{
  "name": "orders-cdc",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql-master",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "***",
    "database.server.id": "5400",
    "database.include.list": "orders_db",
    "table.include.list": "orders_db.orders,orders_db.order_items",
    "topic.prefix": "cdc",
    "schema.history.internal.kafka.topic": "schema-history",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter"
  }
}
// 消费到的事件(Kafka topic: cdc.orders_db.orders)
{
  "op": "u",                                // c=create, u=update, d=delete
  "before": {"id": 1001, "status": "PENDING", "amount": 500},
  "after":  {"id": 1001, "status": "PAID",   "amount": 500},
  "source": {"db": "orders_db", "table": "orders", "lsn": 123456}
}
Flink CDC Connector 把"日志解析"内嵌进 Flink 作业:
  · Flink SQL 直接同步:Source 是 CDC 流,Sink 是目标表
  · 端到端 Exactly-Once(checkpoint + 幂等)
  · 支持动态 schema 演进
  · 一条 SQL 完成"实时入湖/入仓"

典型用法:
  · 实时同步 MySQL → Kafka / Hudi / ClickHouse
  · 双活/跨区同步
  · 实时维表关联
-- Flink SQL:MySQL → Hudi 实时入湖
CREATE TABLE orders_cdc (
  id INT, status STRING, amount DECIMAL(10,2),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'mysql-master',
  'database-name' = 'orders_db',
  'table-name' = 'orders',
  'server-id' = '5400-5404'
);

CREATE TABLE orders_ods (
  id INT, status STRING, amount DECIMAL(10,2),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'hudi',
  'path' = 's3://warehouse/ods/orders'
);

INSERT INTO orders_ods SELECT * FROM orders_cdc;

四、工具选型

工具源支持下游特点
CanalMySQLKafka/MQ/RocketMQ阿里系,轻量
DebeziumMySQL/PG/Mongo/…Kafka Connect 生态跨库、标准化
Flink CDCMySQL/PG/…Flink 任意 SinkSQL 化、Exactly-Once
MaxwellMySQLKafka单表 JSON,极简
Tapdata/流平台多源多目标低代码产品
选型要点:
  · 已上 Kafka Connect → Debezium
  · 主要 MySQL + 阿里云栈 → Canal
  · 要实时入湖/入仓/复杂计算 → Flink CDC
  · 跨多种数据库统一 → Debezium / Flink CDC
  · 团队熟悉度 + 运维成本

五、端到端管道架构

5.1 经典链路

业务库(MySQL/PG)
    │ binlog/WAL
    ▼
CDC (Canal/Debezium/Flink CDC)
    │ 变更事件流
    ▼
Kafka (每个表一个 topic / 分区按主键)
    │ 多路分发(消费者组)
    ├──► Flink → 实时数仓 / ClickHouse / Hudi
    ├──► Redis 缓存刷新
    ├──► Elasticsearch 索引同步
    └──► 审计/归档存储

关键设计:
  · 分区键 = 主键(保证同一行事件有序)
  · 下游消费幂等(upsert)
  · schema 变化通过 registry 协调

5.2 增量 + 全量混合

初次建管道 = 全量快照 + 增量追平:
  1. 先做全量初始同步(如已有备份/历史数据)
  2. 启动 CDC 增量(位点对齐)
  3. 合并去重 → 保证不丢不重

Flink CDC 内置此能力:启动时自动快照 + 增量衔接

六、Schema 演进处理

6.1 DDL 变更的影响

业务库加列/改类型 → CDC 事件结构变化 → 下游 schema 要跟上

三种处理:
  · 统一 schema registry(Avro + Confluent Schema Registry)
  · 宽松 JSON(新增字段容忍)
  · Flink 动态列(重建表 schema)

Debezium 处理:
  · DDL 事件也发到 schema-history topic
  · 下游可从 schema history 重建表结构

6.2 兼容策略

· 只加列(可空/默认值):向后兼容,直接推进
· 改列类型:评估影响面(类型不兼容需重放)
· 删列:下游若引用则报错,需提前治理
· 建议:源库 DDL 走审批 + 通知 CDC 管道

七、生产级管道设计

7.1 可靠性

· 位点持久化(Kafka Connect offset / Flink checkpoint)
· 源库 binlog 保留时长(建议 24-72h,给管道容错窗口)
· 下游幂等(主键 upsert)
· 断点续跑:管道重启从位点继续,无重复无丢失
· 监控:事件延迟(lag)、处理吞吐、失败重试

7.2 事件保序

同一条数据(同主键)的事件必须有序:
  · Kafka 分区按主键 hash → 同 key 同分区 → 有序
  · 下游单分区消费或按 key 聚合处理
  · 跨表(如 order + order_item)天然独立 topic,无需保序

7.3 反压与限流

· 源库 binlog 解析跟不上 → 事件积压 → 延迟上升
  → 扩容 CDC 消费者 / 增大 Kafka 分区
· 下游慢 → 背压传播
  → 用 Kafka 缓冲解耦 + 独立扩容消费组
· 大事务/批量变更 → 事件风暴
  → 下游幂等 + 限流(最大吞吐阈值)

八、监控与运维

8.1 关键指标

# CDC 健康指标
cdc_events_total{table}                  # 事件量
cdc_event_latency_seconds                # 端到端延迟(源库→下游)
kafka_consumer_lag{group, topic}         # 消费积压
cdc_parse_error_total                    # 解析错误
cdc_schema_version{table}                # schema 版本
cdc_binlog_position                     # 位点进度(对比最新)

# 告警:lag 超阈值、parse error 持续、位点停止

8.2 日常运维

· 版本升级:先扩容新版本实例 → 灰度切流
· 位点重置:确认不丢数据才允许
· 表结构变更:走 DDL 审批并同步更新下游
· 故障演练:模拟 binlog 中断 / 消费组故障

九、落地清单与避坑

9.1 Checklist

□ binlog_format=ROW + GTID 开启(MySQL)
□ 授权 CDC 账号(REPLICATION SLAVE/CLIENT 等)
□ 选型:Kafka Connect/Debezium vs Flink CDC vs Canal
□ Kafka topic 设计(每表一 topic,主键分区)
□ 全量+增量衔接(初次建管道)
□ schema registry / 演进策略
□ 位点持久化 + 断点续跑验证
□ 下游幂等 upsert
□ 延迟/积压监控 + 告警
□ binlog 保留时长评估

9.2 常见坑

坑现象对策
binlog 非 ROW无 before/after改 ROW + 换新管道
位点丢失重复/丢失事件持久化 offset
下游不幂等重复数据主键 upsert
schema 变更崩解析失败registry + DDL 审批
binlog 过期无法回追保留 24-72h
事件风暴下游被冲垮限流 + 缓冲
跨表保序误判数据错乱分区键 = 主键

9.3 反模式

· 用轮询 SELECT 代替 CDC(重复拉取、伤源库)→ 应走日志
· 下游直接依赖"事件顺序"做业务逻辑(跨表无序)
· 把 CDC 当无限准确(大事务丢细粒度)→ 配合幂等
· 忽略 schema 演进(加列即崩)

总结:CDC 决策表

环节关键动作
原理binlog/WAL 解析 + 位点管理
工具Canal(MySQL/MQ)/ Debezium(Kafka)/ Flink CDC(SQL)
链路源库 → CDC → Kafka → Flink/缓存/搜索/数仓
可靠位点持久化 + 幂等 + 断点续跑
Schemaregistry + DDL 审批
保序分区键 = 主键
监控延迟 / 积压 / 解析错误

CDC 重新定义了"数据库与其他系统的连接方式"——不再是定时拉取快照,而是实时订阅变更日志。它让数据层第一次拥有了统一的实时通道:缓存、搜索、数仓、审计共享同一条事件流。落地记住四件事:源库日志要开 ROW + 位点要持久化、下游必须幂等、schema 演进要提前治理、延迟与积压要可观测。把 CDC 管道建好,你的数据基础设施就从"批处理的世界"迈入了"流的世界"——实时不再是一句口号,而是默认能力。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. 数据库安全加固与审计实战:权限最小化、加密、脱敏与合规
  2. 数据库容量规划与资源治理:从评估、监控到扩展路径
  3. 数据库字符集、排序规则与乱码实战:utf8mb4、Collation 选择与排查