Flink CDC 从零到一:数据开发工程师的实时数据集成入门指南

写在前面:如果你是一名刚接触实时数据同步的数据开发工程师,一定听过"CDC"这个词。也许你正在为业务库到数仓的 T+1 同步延迟而烦恼,也许你被 Canal + Kafka + 自己写消费者的复杂链路折腾得焦头烂额。这篇博客将带你从零开始,系统学习 Flink CDC 这一当前最流行的开源实时数据集成框架,从原理到实战,从 SQL 到 YAML Pipeline,从单机 Demo 到生产调优,帮你建立完整的知识体系。


目录

  1. CDC 概述
  2. Flink CDC 简介
  3. 核心原理与架构
  4. 环境搭建
  5. Flink CDC SQL 方式(重点)
  6. DataStream API 方式
  7. Sink 集成(重点实战)
  8. 整库同步与数据集成(3.x 核心)
  9. 实时数仓典型架构
  10. 性能调优
  11. 常见问题与踩坑
  12. Flink CDC 3.x 新特性深度解析
  13. 端到端实战项目
  14. 学习路线与资源
  15. 总结

1. CDC 概述

1.1 什么是 CDC

CDC(Change Data Capture,变更数据捕获) 是一种用于捕获数据库中数据变更(INSERT / UPDATE / DELETE)的技术。它能够监测并捕获源数据库的变动记录,将这些变更按发生顺序完整地传递给下游系统。

与传统的批量数据抽取不同,CDC 是事件驱动的——数据一变,即刻捕获,实时性可达毫秒到秒级。这使它成为构建实时数据管道的基石。

1.2 CDC 与传统 ETL 方式的对比

对比维度轮询查询(Polling)触发器(Trigger)双写(Dual Write)日志解析(Log-based CDC)
实时性低(分钟/小时级)高高高(毫秒/秒级)
性能影响频繁查询增加 DB 压力触发器开销大业务侵入强几乎无影响
能否捕获删除困难(软删除可行)能需业务保证能
数据完整性可能丢失中间状态完整易不一致完整
业务侵入无高(建触发器)高(改代码)无
架构复杂度低中高(分布式事务)中
典型工具Sqoop、DataX、Kafka JDBC Source自研存储过程业务代码中实现Canal、Debezium、Flink CDC

日志解析方案在实时性、性能和完整性上优势明显,是目前工业界的主流选择。

1.3 CDC 的两种模式

  • 查询式 CDC(Query-based):通过定时执行 SQL 查询获取变更数据,通常需要时间戳或版本号字段。实现简单但延迟高、无法捕获删除和中间状态。
  • 日志式 CDC(Log-based):直接解析数据库的事务日志(如 MySQL Binlog、PostgreSQL WAL、Oracle Redo Log),实时获取所有变更事件。实现复杂但实时性高、对业务无侵入。

Flink CDC 属于日志式 CDC,底层基于 Debezium 引擎进行日志解析。

1.4 CDC 典型应用场景

  • 实时数仓建设:业务库数据实时同步到 Kafka / Doris / StarRocks / 数据湖
  • 异构数据同步:MySQL → PostgreSQL、Oracle → MySQL 等跨库迁移
  • 缓存更新:数据库变更自动刷新 Redis 缓存,避免 Cache Aside 模式的一致性问题
  • 审计与合规:完整记录每次数据变更的前后值和时间戳
  • 异地多活:跨数据中心的数据复制与同步
  • CQRS(命令查询职责分离):将写模型的变更同步到读模型,支撑高并发查询
  • 数据分发:一份变更数据实时推送到多个下游存储(ES、Redis、OLAP 引擎等)

2. Flink CDC 简介

2.1 Flink CDC 是什么

Flink CDC 是 Apache Flink 社区开源的实时数据集成框架,它提供了一组 Source Connector,能够直接从 MySQL、PostgreSQL、Oracle 等数据库中全量读取历史数据 + 持续读取增量变更,无需借助 Kafka 等中间件中转。

Flink CDC 底层封装了 Debezium 引擎进行日志解析,但在上层做了大量增强:与 Flink 的 Checkpoint 机制深度集成、支持并行快照、提供 SQL 和 YAML 两种声明式 API。

据《Apache Flink CDC 3.2.0 Release Announcement》介绍,截至 3.2 版本,Flink CDC 已支持 MySQL、PostgreSQL、Oracle、SQL Server、MongoDB、DB2、OceanBase、TiDB、Vitess 等数据源。

2.2 发展历史

版本时间核心特性
1.x2020-2021DataStream API;单表读取;基于 Debezium Embedded Engine;SourceFunction 单并发;全量阶段需加锁
2.x2021-2023Flink SQL 支持;增量快照算法(无锁、并行、Checkpoint);整库同步雏形;无锁快照;支持 Oracle / MongoDB / SQL Server 等多数据源;动态加表
3.02023.11架构重构:Pipeline Connector + YAML 声明式整库同步;Schema Evolution 自动同步;Source Coordinator + Split Assigner;分库分表合并;自动缩容
3.12024.05Transform 数据变换(投影/计算/过滤);表合并路由;新增 Kafka / Paimon Pipeline Sink;MySQL tables.exclude 选项
3.22024.09可定制 Schema Evolution(Lenient/TryEvolve 模式 + 事件类型级控制);Transform UDF 支持;复杂路由(一对多广播/模式替换);新增 Elasticsearch Pipeline Sink;K8s 部署模式
3.2.12024.11Bug 修复版本,修复事务泄漏、Schema Evolution 失败、Paimon 重复提交等关键问题

2.3 Flink CDC 核心能力

  • 全量 + 增量一体化:先快照读取存量数据,无缝切换到 Binlog 增量消费
  • 无锁并行快照:增量快照算法无需 FLUSH TABLES WITH READ LOCK,支持多并发读取
  • 断点续传:基于 Flink Checkpoint,chunk 级别容错,故障后从上次完成位置恢复
  • Exactly-Once 语义:与 Flink Checkpoint 协同,保证端到端精确一次
  • 整库同步:YAML 配置一行命令完成整库级数据同步
  • Schema Evolution:上游 DDL 变更自动同步到下游
  • 数据变换与路由:Transform 投影/过滤/计算,Route 支持分库分表合并
  • 丰富的生态:支持写入 Kafka、Doris、StarRocks、Paimon、Iceberg、Hudi、Elasticsearch 等

2.4 支持的数据源

数据源类型支持版本Pipeline SourceSQL Source
MySQL关系型5.6, 5.7, 8.0.x✅✅
PostgreSQL关系型9.6+✅ (3.5+)✅
Oracle关系型11g+✅ (3.6+)✅
SQL Server关系型2012+—✅
MongoDB文档型3.6+—✅
DB2关系型11.5+—✅
OceanBase关系型3.1.x+—✅
TiDB关系型5.1+—✅
Vitess关系型——✅

Pipeline Sink 方面,3.x 已支持 Doris、StarRocks、Kafka、Paimon、Elasticsearch、Iceberg、Hudi(3.6+)、Fluss(3.5+)。

2.5 Flink CDC vs Debezium vs Canal vs Maxwell

对比维度Flink CDCDebeziumCanalMaxwell
定位分布式实时数据集成框架Kafka Connect 上的 CDC 插件MySQL Binlog 订阅组件MySQL Binlog 订阅守护进程
全量同步✅ 内置并行无锁快照✅ 但单线程加锁❌ 需 DataX 等辅助✅ 单表初始化
增量同步✅✅✅✅
并行度快照阶段可水平扩展单 Task(快照可多线程)单线程单线程
锁表❌ 无需锁(3.x 增量快照)⚠️ 全量阶段需锁——
数据源10+ 种10+ 种仅 MySQL仅 MySQL
计算能力✅ Flink SQL/DataStream 强大 ETL❌ 仅捕获+投递❌❌
Schema Evolution✅ 自动同步⚠️ 有限支持❌❌
整库同步✅ YAML 声明式⚠️ 需配置多 Connector❌❌
依赖Flink 集群Kafka Connect 集群ZooKeeper无
高可用Flink HAKafka Connect HAZooKeeper HA无
断点续传Flink CheckpointKafka Offset本地/ZKBinlog 位点
社区活跃度最高(Apache 顶级子项目)高(Red Hat 维护)中(阿里巴巴)低

小结:如果你需要"捕获 + 计算 + 写入"一体化方案,Flink CDC 是最佳选择;如果只需要简单地把 Binlog 投递到 Kafka 且不想引入 Flink,Canal / Maxwell 也够用;Debezium 适合已有 Kafka Connect 生态的团队。


3. 核心原理与架构

3.1 基于数据库日志的解析原理

以 MySQL 为例,Binlog(二进制日志)记录了所有对数据库进行修改的事件(INSERT/UPDATE/DELETE/DDL)。MySQL 主从复制正是依赖 Binlog 实现的。

Flink CDC 伪装成 MySQL 的一个从节点(Slave),向主库发送 COM_BINLOG_DUMP 命令,持续接收 Binlog 事件流,然后通过 Debezium 的解析器将二进制事件解码为结构化的变更记录。

关键前提:

  • Binlog 格式必须为 ROW:STATEMENT 格式记录的是原始 SQL,无法准确获取每行变更前后的数据
  • binlog_row_image=FULL:确保 Binlog 中同时包含变更前和变更后的完整行镜像

PostgreSQL 使用逻辑解码(Logical Decoding)输出 WAL,Oracle 使用 LogMiner 或 XStream,原理类似。

3.2 全量快照 + 增量同步两阶段流程

在这里插入图片描述

两阶段切换的关键:在所有 Chunk 快照读取完成后,Source 会等待一次完整的 Checkpoint 完成(确保所有快照数据已被下游消费),然后从快照结束时记录的 Binlog 位点开始持续消费增量日志,保证数据不丢不重。

3.3 Flink CDC 2.x 架构

2.x 版本基于 Flink 的 FLIP-27 Source API 重构,引入了增量快照框架:

  • MySqlSourceEnumerator(运行在 JobManager):负责将表按主键切分为 Chunk,并将 Chunk 分配给 Reader
  • MySqlSourceReader(运行在 TaskManager):并行读取分配到的 Chunk 数据
  • 快照阶段多 Reader 并行,增量阶段单 Reader 读取 Binlog 保证顺序
  • 底层仍使用 Debezium Embedded Engine 解析 Binlog

但 2.x 仍有局限:只面向单表 SQL Source,缺乏端到端整库同步能力,Schema 变更处理薄弱。

3.4 Flink CDC 3.x 架构重构

3.x 是一次重大架构升级,定位从"CDC Connector 集合"升级为端到端实时数据集成框架。

在这里插入图片描述

3.4.1 分片(Chunk)并行快照

表按 Chunk Key(默认为主键的第一列)切分为多个 Chunk,每个 Chunk 是一个独立的数据分片。Enumerator 将 Chunk 分发给多个 Reader 并行读取,水平扩展快照吞吐量。

3.2 版本起支持将任意列指定为 Chunk Key(scan.incremental.snapshot.chunk.key-column),不再局限于主键列,但非主键列可能导致查询性能下降。

3.4.2 增量快照算法(Offset Signal Algorithm)

Flink CDC 的增量快照算法受 Netflix DBLog 论文启发,核心步骤如下(以单个 Chunk 为例):

  1. 记录当前 Binlog 位点为 LOW 位点
  2. 执行 SELECT * FROM table WHERE pk > chunk_low AND pk <= chunk_high 读取该 Chunk 的快照数据并缓存
  3. 记录当前 Binlog 位点为 HIGH 位点
  4. 读取 LOW 到 HIGH 之间属于该 Chunk 的 Binlog 变更记录
  5. 将 Binlog 变更 Upsert 到缓存的快照数据中(以最新状态为准)
  6. 将最终结果全部作为 INSERT 事件发出
  7. HIGH 位点之后的 Binlog 记录在增量阶段由单 Reader 继续读取

该算法的三大优势:

  • 无锁:不需要 FLUSH TABLES WITH READ LOCK
  • 支持断点续传:Chunk 粒度 Checkpoint,故障后从最后完成的 Chunk 恢复
  • 并发读取:多个 Chunk 可由不同 Reader 并行处理
3.4.3 Source Coordinator + Split Assigner

3.x 在 Pipeline 层面引入了统一的 Source Coordinator,负责:

  • 管理所有表的 Chunk 切分与分配
  • 协调快照阶段到增量阶段的切换
  • 支持动态新增表(运行中添加监控表,自动触发该表快照并接入 Binlog 流)
  • 自动缩容(快照完成后释放多余 Reader 资源)
3.4.4 数据变更 Event 模型

3.x 定义了统一的事件模型:

  • DataChangeEvent:数据变更事件,包含 before(变更前行数据)、after(变更后行数据)、op 操作类型(INSERT/UPDATE/DELETE)、元数据(时间戳、表名等)
  • SchemaChangeEvent:结构变更事件,包括 CreateTableEvent、AddColumnEvent、DropColumnEvent、RenameColumnEvent、AlterColumnTypeEvent、AlterTableCommentEvent、AlterColumnCommentEvent 等

Schema Operator 在收到 SchemaChangeEvent 时,会先向下游发送 FlushEvent 等待 Sink 将旧 Schema 的数据全部落盘,再通过 MetadataApplier 将 DDL 应用到下游,保证数据与 Schema 的一致性。

3.5 水位线与 Checkpoint 协同

Flink CDC 与 Flink 的 Checkpoint 机制深度集成:

  • 快照阶段:以 Chunk 为粒度做 Checkpoint,每个 Chunk 读取完成后记录进度
  • 增量阶段:以行为粒度记录 Binlog 消费位点到 State
  • 两阶段切换屏障:所有 Chunk 完成后,等待一次成功的 Checkpoint,确保快照数据全部被下游消费,再切换到 Binlog 读取
  • 故障恢复:从最近一次成功 Checkpoint 的 State 恢复,重新分配未完成的 Chunk,重置 Binlog 位点,保证 Exactly-Once

4. 环境搭建

4.1 前置条件

组件版本要求
JDK11+(3.6+ 强制 JDK 11)
Flink1.17+(3.x),推荐 1.18 / 1.20
MySQL5.6 / 5.7 / 8.0.x
内存开发环境至少 4GB,生产按需

4.2 MySQL Binlog 配置

编辑 MySQL 配置文件 my.cnf(Linux)或 my.ini(Windows):

[mysqld]
# 开启 Binlog
log-bin=mysql-bin
# Binlog 格式必须为 ROW
binlog_format=ROW
# 记录完整行镜像(前后镜像都有)
binlog_row_image=FULL
# server-id 必须唯一
server-id=1
# Binlog 过期时间(天),根据磁盘容量调整
expire_logs_days=7
#(可选)开启 GTID,方便主从切换和位点管理
gtid_mode=ON
enforce_gtid_consistency=ON

重启 MySQL 后验证:

-- 检查 Binlog 是否开启
SHOW VARIABLES LIKE 'log_bin';
-- 应为 ROW
SHOW VARIABLES LIKE 'binlog_format';
-- 应为 FULL
SHOW VARIABLES LIKE 'binlog_row_image';

创建 CDC 专用账号并授权:

CREATE USER 'flinkcdc'@'%' IDENTIFIED BY 'FlinkCDC@123';
GRANT SELECT, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flinkcdc'@'%';
FLUSH PRIVILEGES;

注意:3.x 增量快照算法无需 RELOAD 权限(旧版本加锁快照需要),权限更小更安全。

4.3 Flink Standalone 本地环境

# 下载 Flink(以 1.18.1 为例)
wget https://archive.apache.org/dist/flink/flink-1.18.1/flink-1.18.1-bin-scala_2.12.tgz
tar -xzf flink-1.18.1-bin-scala_2.12.tgz
cd flink-1.18.1

# 启动集群
./bin/start-cluster.sh

# 访问 http://localhost:8081 查看 Web UI

4.4 Flink CDC 依赖下载

从 Flink CDC Releases 下载对应版本的 Connector JAR:

# 以 3.2.1 为例
cd $FLINK_HOME/lib

# SQL Connector(Flink SQL 方式使用)
wget https://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-mysql-cdc/3.2.1/flink-sql-connector-mysql-cdc-3.2.1.jar

# Pipeline Connector(YAML 整库同步方式使用)
wget https://repo1.maven.org/maven2/org/apache/flink/flink-cdc-pipeline-connector-mysql/3.2.1/flink-cdc-pipeline-connector-mysql-3.2.1.jar
wget https://repo1.maven.org/maven2/org/apache/flink/flink-cdc-pipeline-connector-doris/3.2.1/flink-cdc-pipeline-connector-doris-3.2.1.jar

重要:MySQL / Oracle / OceanBase / DB2 连接器因许可证原因不包含 JDBC 驱动,需手动将驱动 JAR 放入 $FLINK_HOME/lib/。MySQL 推荐使用 mysql-connector-java-8.0.27.jar。

4.5 Flink SQL Client 快速验证

# 启动 SQL Client
./bin/sql-client.sh

在 SQL Client 中执行:

-- 创建 MySQL CDC 源表
CREATE TABLE mysql_users (
    id INT,
    name STRING,
    email STRING,
    created_at TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'demo',
    'table-name' = 'users',
    'server-id' = '5400-5404'
);

-- 查询(先全量读取已有数据,然后持续输出增量变更)
SELECT * FROM mysql_users;

在 MySQL 中插入或更新一条数据,SQL Client 中应该能实时看到变更。

4.6 Docker Compose 一键启动

创建 docker-compose.yml:

version: '3.8'
services:
  mysql:
    image: mysql:8.0
    container_name: mysql-cdc
    environment:
      MYSQL_ROOT_PASSWORD: root
      MYSQL_DATABASE: demo
    ports:
      - "3306:3306"
    command:
      - --server-id=1
      - --log-bin=mysql-bin
      - --binlog_format=ROW
      - --binlog_row_image=FULL
    volumes:
      - ./init.sql:/docker-entrypoint-initdb.d/init.sql

  flink-jobmanager:
    image: flink:1.18.1-java11
    container_name: flink-jm
    ports:
      - "8081:8081"
    command: jobmanager
    environment:
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: flink-jobmanager
    volumes:
      - ./lib:/opt/flink/lib

  flink-taskmanager:
    image: flink:1.18.1-java11
    container_name: flink-tm
    depends_on:
      - flink-jobmanager
    command: taskmanager
    environment:
      FLINK_PROPERTIES: |
        jobmanager.rpc.address: flink-jobmanager
        taskmanager.numberOfTaskSlots: 4
    volumes:
      - ./lib:/opt/flink/lib

创建 init.sql:

CREATE DATABASE IF NOT EXISTS demo;
USE demo;
CREATE TABLE users (
    id INT AUTO_INCREMENT PRIMARY KEY,
    name VARCHAR(100),
    email VARCHAR(200),
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
INSERT INTO users (name, email) VALUES ('Alice', 'alice@example.com'), ('Bob', 'bob@example.com');

启动:

docker-compose up -d

4.7 常见启动问题排查

问题原因解决方案
Cannot read binlogBinlog 未开启或格式不对确认 log_bin=ON、binlog_format=ROW
Access denied账号权限不足授予 REPLICATION SLAVE / CLIENT 权限
server-id conflict多个 CDC 任务使用了相同 server-id为每个任务配置不同的 server-id 范围
ClassNotFoundException: com.mysql.cj.jdbc.Driver缺少 MySQL 驱动将 mysql-connector-java JAR 放入 lib/
Timezone mismatch未配置时区添加 'server-time-zone' = 'Asia/Shanghai'
OutOfMemoryError快照数据量过大增大 TaskManager 内存或调小 chunk.size

5. Flink CDC SQL 方式(重点)

Flink SQL 是最常用、最易上手的方式。只需通过 DDL 定义源表和目标表,然后用 INSERT INTO 语句启动同步。

5.1 MySQL CDC Source DDL 详解

CREATE TABLE mysql_source (
    id INT,
    name STRING,
    age INT,
    email STRING,
    created_at TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    -- ========== 必选参数 ==========
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'demo',
    'table-name' = 'users',

    -- ========== 连接与身份 ==========
    'server-id' = '5400-5404',
    'server-time-zone' = 'Asia/Shanghai',

    -- ========== 快照配置 ==========
    'scan.incremental.snapshot.enabled' = 'true',
    'scan.incremental.snapshot.chunk.size' = '8096',
    'scan.incremental.snapshot.chunk.key-column' = 'id',
    'scan.snapshot.fetch.size' = '1024',

    -- ========== 启动模式 ==========
    'scan.startup.mode' = 'initial',

    -- ========== 连接池 ==========
    'connection.pool.size' = '20',
    'connect.timeout' = '30s',
    'heartbeat.interval' = '30s'
);

核心 WITH 参数速查表:

参数必填默认值说明
connector✅—固定为 mysql-cdc
hostname✅—MySQL 主机地址
port❌3306MySQL 端口
username✅—数据库用户名
password✅—数据库密码
database-name✅—数据库名,支持正则
table-name✅—表名,支持正则
server-id❌随机 5400-6400MySQL 从节点 ID,并行时需配置范围
server-time-zone❌会话时区数据库会话时区,如 Asia/Shanghai
scan.incremental.snapshot.enabled❌true是否启用增量快照(并行无锁)
scan.incremental.snapshot.chunk.size❌8096每个 Chunk 的行数
scan.incremental.snapshot.chunk.key-column❌主键第一列Chunk 切分列
scan.snapshot.fetch.size❌1024每次轮询获取的最大行数
scan.startup.mode❌initial启动模式(见 5.2)
connection.pool.size❌20JDBC 连接池大小
heartbeat.interval❌30s心跳间隔,防止连接超时

注意:PRIMARY KEY 声明是必须的(用于 Upsert 和 Chunk 切分),Flink CDC 中使用 NOT ENFORCED 表示不强制校验。

5.2 启动模式

模式说明适用场景
initial(默认)首次启动先全量快照,再从快照结束位点读取增量首次同步,需要历史数据
latest-offset跳过快照,仅从 Binlog 末尾开始读取只关心启动后的新变更
earliest-offset跳过快照,从 Binlog 最早位点开始读取需要回放所有历史 Binlog
specific-offset从指定 Binlog 文件和位置或 GTID 集合开始精确恢复位点
timestamp从指定时间戳对应的 Binlog 位点开始按时间点恢复

specific-offset 模式额外参数:

'scan.startup.mode' = 'specific-offset',
'scan.startup.specific-offset.file' = 'mysql-bin.000003',
'scan.startup.specific-offset.pos' = '4',
-- 或使用 GTID
'scan.startup.specific-offset.gtid-set' = '24DA167-0C0C-11E8-8442-00059A3C7B00:1-19'

timestamp 模式:

'scan.startup.mode' = 'timestamp',
'scan.startup.timestamp-millis' = '1700000000000'

5.3 全量 + 增量同步完整 SQL 示例

-- ============================================
-- 1. 创建 MySQL CDC 源表
-- ============================================
CREATE TABLE orders_source (
    id BIGINT,
    user_id BIGINT,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'ecommerce',
    'table-name' = 'orders',
    'server-id' = '5400-5404',
    'server-time-zone' = 'Asia/Shanghai',
    'scan.startup.mode' = 'initial'
);

-- ============================================
-- 2. 创建目标表(以 Doris 为例)
-- ============================================
CREATE TABLE orders_sink (
    id BIGINT,
    user_id BIGINT,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = '127.0.0.1:8030',
    'table.identifier' = 'ecommerce.orders',
    'username' = 'root',
    'password' = '',
    'sink.label-prefix' = 'orders_sync',
    'sink.properties.format' = 'json',
    'sink.properties.read_json_by_line' = 'true'
);

-- ============================================
-- 3. 启动同步
-- ============================================
INSERT INTO orders_sink
SELECT * FROM orders_source;

提交后,Flink 会先全量读取 orders 表中的存量数据写入 Doris,然后持续监听 Binlog 将增量变更实时写入。

5.4 增量同步(只从 Binlog 开始)

如果历史数据已经通过其他方式(如 DataX)同步完毕,只需增量同步:

CREATE TABLE orders_cdc (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'ecommerce',
    'table-name' = 'orders',
    'server-id' = '5400',
    'scan.startup.mode' = 'latest-offset'
);

5.5 表名正则匹配多表

table-name 支持正则表达式,可以同时监听多张表:

CREATE TABLE multi_table_source (
    id INT,
    name STRING,
    -- 元数据列:标识数据来自哪张表
    table_name STRING METADATA FROM 'table_name' VIRTUAL,
    db_name STRING METADATA FROM 'database_name' VIRTUAL,
    op_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL,
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'ecommerce',
    -- 匹配 order_0, order_1, order_2 ...
    'table-name' = 'order_[0-9]+',
    'server-id' = '5400-5408'
);

注意:使用正则匹配多表时,所有匹配的表必须具有相同的 Schema 结构(列名和类型一致)。如果表结构不同,应使用 3.x 的 YAML Pipeline 整库同步方式。

5.6 整库同步(3.x):YAML Pipeline 方式

3.x 引入了 YAML 声明式整库同步,这是对 SQL 方式的重大升级。详见第 8 章。这里先给出一个最简示例:

source:
  type: mysql
  hostname: 127.0.0.1
  port: 3306
  username: flinkcdc
  password: FlinkCDC@123
  tables: ecommerce.\.*
  server-id: 5401-5404

sink:
  type: doris
  fenodes: 127.0.0.1:8030
  username: root
  password: ""

pipeline:
  name: MySQL-to-Doris-Sync
  parallelism: 4
bash flink-cdc.sh mysql-to-doris.yaml

5.7 Schema Evolution(3.x)

3.x 支持在 YAML Pipeline 中自动同步上游 Schema 变更。当上游 MySQL 执行 DDL 时,Flink CDC 会捕获 SchemaChangeEvent 并自动在下游执行对应 DDL。

source:
  type: mysql
  hostname: 127.0.0.1
  # ...
  schema-change.enabled: true

sink:
  type: doris
  # ...
  # Schema Evolution 行为配置
  include.schema.changes: [create_table, add_column, alter_column_type]
  exclude.schema.changes: [drop_column]

pipeline:
  name: Schema-Evolution-Demo
  parallelism: 2

Schema Evolution 支持以下行为模式(通过 pipeline.schema.change.behavior 配置):

模式说明
lenient(3.2+ 默认)宽容模式:忽略删列操作,新增列追加到末尾,故障恢复时无损
try_evolve尝试演化:尽力应用 Schema 变更,失败则容忍(可能丢数据)
evolve严格演化:应用所有 Schema 变更,失败则终止任务
ignore忽略:忽略所有 Schema 变更,仅传播未变更列的数据
exception异常模式:禁止任何 Schema 变更,发生即终止任务

5.8 元数据列

Flink CDC 支持通过 METADATA FROM 语法读取变更事件的元数据:

CREATE TABLE orders_with_meta (
    id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    -- 元数据列
    op_ts TIMESTAMP_LTZ(3) METADATA FROM 'op_ts' VIRTUAL,        -- 变更时间
    table_name STRING METADATA FROM 'table_name' VIRTUAL,         -- 表名
    database_name STRING METADATA FROM 'database_name' VIRTUAL,  -- 库名
    op STRING METADATA FROM 'op' VIRTUAL,                        -- 操作类型(INSERT/UPDATE/DELETE)
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '127.0.0.1',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'ecommerce',
    'table-name' = 'orders'
);

常用元数据键:

  • op_ts:变更发生的时间戳
  • table_name:源表名
  • database_name:源库名
  • op:操作类型(INSERT / UPDATE / DELETE)
  • meta.table_name:多表场景中的完整表标识

6. DataStream API 方式

虽然 SQL 方式最常用,但在需要高度定制化逻辑(如自定义反序列化、复杂的多表合并、精细的并行度控制)时,DataStream API 提供了更强的灵活性。

6.1 SourceFunction 方式(2.x 旧 API)

2.x 时代使用 MySqlSource.Builder 构建 Source(注意:这是 FLIP-27 的新 Source API,不是旧的 SourceFunction):

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;

public class OldApiExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env =
                StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(5000);

        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
                .hostname("127.0.0.1")
                .port(3306)
                .databaseList("ecommerce")
                .tableList("ecommerce.orders")
                .username("flinkcdc")
                .password("FlinkCDC@123")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .serverId("5400-5404")
                .build();

        DataStream<String> stream = env.fromSource(
                mySqlSource,
                WatermarkStrategy.noWatermarks(),
                "MySQL CDC Source"
        );

        stream.print();
        env.execute("Flink CDC Old API Example");
    }
}

注意:JsonDebeziumDeserializationSchema 输出的是 Debezium 标准 JSON 格式字符串,包含 before / after / source / op / ts_ms 字段。

6.2 3.x 新 Source API

3.x 的 DataStream API 保持了与 2.x 类似的 Source 构建方式,主要变化在于 Pipeline 层面的 API。对于纯 Source 读取,API 基本一致:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class Cdc3Example {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env =
                StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(10000); // 10秒一次 Checkpoint
        env.setParallelism(4);

        MySqlSource<String> source = MySqlSource.<String>builder()
                .hostname("127.0.0.1")
                .port(3306)
                .username("flinkcdc")
                .password("FlinkCDC@123")
                .databaseList("ecommerce")
                .tableList("ecommerce.orders", "ecommerce.users")
                .serverId("5400-5408")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .startupOptions(StartupOptions.initial())
                // 增量快照配置
                .splitSize(8096)
                .fetchSize(1024)
                .serverTimeZone("Asia/Shanghai")
                .build();

        env.fromSource(source, WatermarkStrategy.noWatermarks(), "MySQL CDC")
           .setParallelism(4)
           .print()
           .setParallelism(1);

        env.execute("Flink CDC 3.x DataStream Example");
    }
}

Maven 依赖:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-mysql-cdc</artifactId>
    <version>3.2.1</version>
</dependency>

6.3 全量 + 增量 Java 代码示例

以下示例演示从 MySQL 读取订单变更,做简单的 ETL 处理后写入 Kafka:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class FullIncrementalSync {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env =
                StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(10000);
        env.setParallelism(4);

        // 1. MySQL CDC Source
        MySqlSource<String> source = MySqlSource.<String>builder()
                .hostname("127.0.0.1")
                .port(3306)
                .username("flinkcdc")
                .password("FlinkCDC@123")
                .databaseList("ecommerce")
                .tableList("ecommerce.orders")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .serverId("5400-5408")
                .startupOptions(StartupOptions.initial())
                .serverTimeZone("Asia/Shanghai")
                .build();

        // 2. Kafka Sink
        KafkaSink<String> kafkaSink = KafkaSink.<String>builder()
                .setBootstrapServers("127.0.0.1:9092")
                .setRecordSerializer(
                        KafkaRecordSerializationSchema.<String>builder()
                                .setTopic("ods_orders_binlog")
                                .setValueSerializationSchema(new SimpleStringSchema())
                                .build()
                )
                .build();

        // 3. 读取 CDC 数据并写入 Kafka
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "MySQL CDC")
           .sinkTo(kafkaSink);

        env.execute("MySQL CDC -> Kafka");
    }
}

6.4 自定义反序列化(DebeziumDeserializationSchema)

默认的 JsonDebeziumDeserializationSchema 输出 Debezium 原生 JSON,可能不符合下游需求。可以实现 DebeziumDeserializationSchema 自定义输出格式:

import org.apache.flink.cdc.debezium.DebeziumDeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.util.Collector;
import org.apache.kafka.connect.data.Field;
import org.apache.kafka.connect.data.Struct;
import org.apache.kafka.connect.source.SourceRecord;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;

public class CustomDebeziumDeserializer
        implements DebeziumDeserializationSchema<String> {

    private static final ObjectMapper MAPPER = new ObjectMapper();

    @Override
    public void deserialize(SourceRecord record, Collector<String> out)
            throws Exception {
        Struct value = (Struct) record.value();
        String op = value.getString("op"); // c=insert, u=update, d=delete
        Struct after = value.getStruct("after");
        Struct before = value.getStruct("before");
        Struct source = value.getStruct("source");

        ObjectNode result = MAPPER.createObjectNode();
        result.put("table", source.getString("table"));
        result.put("database", source.getString("db"));
        result.put("op", mapOp(op));
        result.put("ts_ms", value.getInt64("ts_ms"));

        // 解析 after 数据
        if (after != null) {
            ObjectNode afterNode = MAPPER.createObjectNode();
            for (Field field : after.schema().fields()) {
                Object fieldValue = after.get(field);
                afterNode.put(field.name(),
                        fieldValue == null ? null : fieldValue.toString());
            }
            result.set("after", afterNode);
        }

        // 解析 before 数据(UPDATE/DELETE 时有值)
        if (before != null) {
            ObjectNode beforeNode = MAPPER.createObjectNode();
            for (Field field : before.schema().fields()) {
                Object fieldValue = before.get(field);
                beforeNode.put(field.name(),
                        fieldValue == null ? null : fieldValue.toString());
            }
            result.set("before", beforeNode);
        }

        out.collect(MAPPER.writeValueAsString(result));
    }

    private String mapOp(String op) {
        switch (op) {
            case "c": return "INSERT";
            case "u": return "UPDATE";
            case "d": return "DELETE";
            default:  return op;
        }
    }

    @Override
    public TypeInformation<String> getProducedType() {
        return TypeInformation.of(String.class);
    }
}

使用自定义反序列化器:

MySqlSource<String> source = MySqlSource.<String>builder()
        .hostname("127.0.0.1")
        .port(3306)
        .username("flinkcdc")
        .password("FlinkCDC@123")
        .databaseList("ecommerce")
        .tableList("ecommerce.orders")
        .deserializer(new CustomDebeziumDeserializer())
        .serverId("5400-5408")
        .build();

6.5 多表合并与路由

DataStream API 可以通过 tableList 指定多张表,在反序列化器中根据表名分发处理:

MySqlSource<String> source = MySqlSource.<String>builder()
        // ...
        .tableList(
            "ecommerce.orders",
            "ecommerce.order_items",
            "ecommerce.users"
        )
        .deserializer(new CustomDebeziumDeserializer())
        .build();

DataStream<String> stream = env.fromSource(source,
        WatermarkStrategy.noWatermarks(), "MySQL CDC");

// 根据表名分流
stream.filter(json -> json.contains("\"table\":\"orders\""))
      .sinkTo(ordersKafkaSink);

stream.filter(json -> json.contains("\"table\":\"users\""))
      .sinkTo(usersKafkaSink);

提示:对于多表合并和路由场景,3.x 的 YAML Pipeline Route 功能更加简洁,详见第 8 章。

6.6 并行度与分片配置

MySqlSource<String> source = MySqlSource.<String>builder()
        // ...
        .serverId("5400-5410")  // server-id 范围必须大于并行度
        .splitSize(10000)       // 每个 Chunk 10000 行
        .fetchSize(2000)        // 每次 fetch 2000 行
        .build();

env.setParallelism(8);  // 快照阶段 8 并发读取

关键规则:

  • server-id 范围大小必须大于 Source 的并行度(每个 Reader 需要一个唯一 server-id)
  • 并行度通过 env.setParallelism() 或 SET 'parallelism.default' 设置
  • 增量阶段(Binlog 读取)始终为单并发,以保证事件的全局顺序

7. Sink 集成(重点实战)

Flink CDC 的数据可以写入多种下游系统。以下介绍常用 Sink 的配置方法和注意事项。

7.1 Kafka Sink

Kafka 是最常用的 CDC 数据落地目标,通常作为 ODS 层的实时数据接入层。

-- Kafka Sink 表(使用 Canal-JSON 格式)
CREATE TABLE kafka_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'kafka',
    'topic' = 'ods_orders_binlog',
    'properties.bootstrap.servers' = '127.0.0.1:9092',
    'properties.group.id' = 'flink-cdc-group',
    -- Canal-JSON 格式:包含 UPDATE 前后镜像
    'format' = 'canal-json',
    -- 或使用 Debezium-JSON
    -- 'format' = 'debezium-json',
    'scan.startup.mode' = 'earliest-offset'
);

-- 启动同步
INSERT INTO kafka_orders
SELECT * FROM orders_source;

格式选择建议:

格式特点适用场景
canal-json包含 old 字段(变更前数据),国内生态友好国内数仓体系,下游 Flink/Spark 消费
debezium-jsonDebezium 标准格式,包含 before/after/source国际化生态,与 Debezium 兼容
avro-confluentSchema Registry + Avro,紧凑高效大规模数据,Schema 治理严格
raw仅原始字节自定义格式透传

Exactly-Once:Kafka Sink 支持两阶段提交(2PC),配合 Flink Checkpoint 可实现端到端精确一次。需设置:

'sink.delivery-guarantee' = 'exactly-once',
'sink.transactional-id-prefix' = 'cdc-txn-'

7.2 Doris Sink

Doris 是国内非常流行的实时分析型数据库,Flink Doris Connector 通过 Stream Load 实现高效写入。

CREATE TABLE doris_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = '127.0.0.1:8030',
    'table.identifier' = 'ecommerce.orders',
    'username' = 'root',
    'password' = '',
    'sink.label-prefix' = 'orders_cdc',
    -- Stream Load 属性
    'sink.properties.format' = 'json',
    'sink.properties.read_json_by_line' = 'true',
    -- 写入参数
    'sink.buffer-flush.max-rows' = '100000',
    'sink.buffer-flush.max-bytes' = '104857600',
    'sink.buffer-flush.interval' = '10s',
    'sink.max-retries' = '3'
);

INSERT INTO doris_orders
SELECT * FROM orders_source;

幂等说明:Doris Unique Key 模型通过主键实现 Upsert,Stream Load 的 Label 机制保证请求幂等,可达到 Exactly-Once 效果。

7.3 StarRocks Sink

StarRocks Sink 配置与 Doris 类似:

CREATE TABLE sr_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'starrocks',
    'jdbc-url' = 'jdbc:mysql://127.0.0.1:9030',
    'load-url' = '127.0.0.1:8030',
    'database-name' = 'ecommerce',
    'table-name' = 'orders',
    'username' = 'root',
    'password' = '',
    'sink.properties.format' = 'json',
    'sink.properties.strip_outer_array' = 'true',
    'sink.buffer-flush.interval-ms' = '10000',
    'sink.buffer-flush.max-bytes' = '104857600'
);

INSERT INTO sr_orders
SELECT * FROM orders_source;

StarRocks 同样通过主键模型实现 Upsert 幂等写入。3.2 版本起 StarRocks Pipeline Sink 支持列重命名和列类型变更 DDL。

7.4 Paimon Sink

Apache Paimon 是流批一体的湖存储格式,Flink CDC 3.1 起提供原生 Pipeline Sink 支持。

YAML Pipeline 方式(推荐):

source:
  type: mysql
  hostname: 127.0.0.1
  port: 3306
  username: flinkcdc
  password: FlinkCDC@123
  tables: ecommerce.\.*
  server-id: 5401-5404

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: filesystem
  catalog.properties.warehouse: /data/paimon/warehouse

pipeline:
  name: MySQL-to-Paimon
  parallelism: 4

Flink SQL 方式:

-- 创建 Paimon Catalog
CREATE CATALOG paimon_catalog WITH (
    'type' = 'paimon',
    'warehouse' = '/data/paimon/warehouse'
);

-- 在 Paimon 中创建表
CREATE TABLE paimon_catalog.ecommerce.orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'bucket' = '4',
    'changelog-producer' = 'input'  -- 直接使用 CDC 输入作为 changelog
);

-- 写入
INSERT INTO paimon_catalog.ecommerce.orders
SELECT * FROM orders_source;

Exactly-Once:Paimon 基于 Snapshot 和两阶段提交,原生支持 Exactly-Once。changelog-producer = 'input' 表示直接保留 CDC 的 changelog 语义(INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE)。

7.5 Iceberg Sink

CREATE CATALOG iceberg_catalog WITH (
    'type' = 'iceberg',
    'catalog-type' = 'hive',
    'uri' = 'thrift://127.0.0.1:9083',
    'warehouse' = 'hdfs:///user/hive/warehouse'
);

CREATE TABLE iceberg_catalog.ecommerce.orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'format-version' = '2',
    'write.upsert.enabled' = 'true'
);

INSERT INTO iceberg_catalog.ecommerce.orders
SELECT * FROM orders_source;

Iceberg v2 表格式支持行级删除(Row-Level Delete),通过 Upsert 实现 CDC 数据写入。Flink CDC 3.x Pipeline 也提供了 Iceberg Sink Connector。

7.6 Hudi Sink

CREATE TABLE hudi_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'hudi',
    'path' = 'hdfs:///user/hudi/orders',
    'table.type' = 'MERGE_ON_READ',
    'hoodie.datasource.write.recordkey.field' = 'id',
    'hoodie.datasource.write.precombine.field' = 'order_time',
    'write.operation' = 'upsert',
    'write.tasks' = '4'
);

INSERT INTO hudi_orders
SELECT * FROM orders_source;

Hudi MOR 表通过 Compaction 合并增量数据,适合高频更新场景。Flink CDC 3.6 新增了 Hudi Pipeline Sink Connector。

7.7 JDBC Sink(MySQL/PostgreSQL)

CREATE TABLE jdbc_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:mysql://127.0.0.1:3306/ecommerce_dwd',
    'table-name' = 'orders',
    'username' = 'root',
    'password' = 'root',
    'driver' = 'com.mysql.cj.jdbc.Driver',
    -- Upsert 模式
    'sink.buffer-flush.max-rows' = '1000',
    'sink.buffer-flush.interval' = '2s',
    'sink.max-retries' = '3'
);

注意:Flink JDBC Sink 在处理 UPDATE / DELETE 时需要目标表有主键约束。对于 MySQL,会自动生成 INSERT ... ON DUPLICATE KEY UPDATE 语句实现 Upsert。

7.8 Redis Sink

Redis 通常作为缓存层或实时指标存储:

CREATE TABLE redis_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),

    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'redis',
    'redis-mode' = 'single',
    'host' = '127.0.0.1',
    'port' = '6379',
    'password' = '',
    'command' = 'SET',
    'ttl' = '86400'
);

Redis 写入是幂等的(SET 覆盖写),但无法表达 DELETE 语义,需要通过 TTL 或特殊值处理。

7.9 Elasticsearch Sink

CREATE TABLE es_orders (
    id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'elasticsearch-7',
    'hosts' = 'http://127.0.0.1:9200',
    'index' = 'orders',
    'document-id.key-delimiter' = '_',
    'sink.bulk-flush.max-actions' = '1000',
    'sink.bulk-flush.interval' = '5s',
    'sink.bulk-flush.backoff.strategy' = 'EXPONENTIAL',
    'sink.bulk-flush.backoff.max-retries' = '3'
);

INSERT INTO es_orders
SELECT * FROM orders_source;

Flink CDC 3.2 新增了 Elasticsearch Pipeline Sink Connector,支持在 ES 6.8 / 7.10 / 8.12 上验证通过。ES 的文档 ID 为主键,写入天然幂等。

7.10 各 Sink 语义对比

Sink写入方式Upsert 支持Delete 支持Exactly-Once
KafkaProducer 发送❌(追加)✅(消息中携带)✅(2PC 事务)
DorisStream Load✅(Unique Key)✅✅(Label 幂等)
StarRocksStream Load✅(Primary Key)✅✅(Label 幂等)
PaimonSnapshot Commit✅✅✅
IcebergSnapshot Commit✅(v2)✅(v2)✅
HudiCoW / MoR✅✅✅
JDBCINSERT/UPSERT✅(主键)⚠️(需配置)⚠️(需幂等键)
RedisSET/HSET✅⚠️(特殊处理)❌(At-Least-Once)
ElasticsearchBulk API✅(doc ID)✅⚠️(At-Least-Once)

8. 整库同步与数据集成(3.x 核心)

3.x 最大的变革是引入了 YAML Pipeline 作为端到端整库同步的声明式方式,将数据集成从"写代码"简化为"写配置"。

8.1 整库同步架构

在这里插入图片描述

四大核心组件:

  • Source:配置数据源连接信息和监控的表
  • Transform:数据变换(投影列、过滤行、计算列、UDF)
  • Route:路由规则(表名映射、分库分表合并、一对多广播)
  • Sink:目标端配置

8.2 YAML 整库同步配置完整示例

# ============================================
# MySQL 整库同步到 Doris 完整配置
# ============================================

pipeline:
  name: MySQL-to-Doris-Full-Sync
  parallelism: 4
  schema.change.behavior: lenient

source:
  type: mysql
  name: MySQL Source
  hostname: 127.0.0.1
  port: 3306
  username: flinkcdc
  password: FlinkCDC@123
  # 监控 ecommerce 库下所有表
  tables: ecommerce.\.*
  # 排除不需要的表
  tables.exclude: ecommerce.temp_.*, ecommerce.log_.*
  server-id: 5401-5408
  server-time-zone: Asia/Shanghai
  scan.startup.mode: initial
  scan.incremental.snapshot.chunk.size: 8096
  schema-change.enabled: true

sink:
  type: doris
  name: Doris Sink
  fenodes: 127.0.0.1:8030
  username: root
  password: ""
  table.create.properties.light_schema_change: true
  table.create.properties.replication_num: 1
  sink.buffer-flush.interval: 1s
  sink.buffer-flush.max-rows: 10000
  sink.buffer-flush.max-bytes: 10485760

# ============================================
# Transform:数据变换
# ============================================
transform:
  - source-table: ecommerce.orders
    projection: id, user_id, product_id, quantity, amount, status, order_time,
                UPPER(status) AS status_upper,
                amount * quantity AS total_amount
    filter: status IS NOT NULL AND amount > 0

  - source-table: ecommerce.users
    projection: id, name, email, created_at,
                CONCAT(name, '_', id) AS user_key
    filter: id > 0

# ============================================
# Route:表名路由与分库分表合并
# ============================================
route:
  # 分表合并:order_0, order_1, ... → merged_orders
  - source-table: ecommerce.order_[0-9]+
    sink-table: ecommerce.merged_orders

  # 模式替换:所有源表名加 ods_ 前缀
  - source-table: ecommerce.\.*
    sink-table: ods_ecommerce.<>
    replace-symbol: "<>"

  # 一对多广播:将重要配置表同步到多个目标
  - source-table: ecommerce.config
    sink-table:
      - ecommerce.config_backup
      - ecommerce.config_archive

8.3 表名/库名路由规则

基础路由:

route:
  - source-table: ecommerce.orders
    sink-table: dwd.dwd_orders

正则匹配 + 模式替换(3.2+):

route:
  # 所有 ecommerce 库的表同步到 dwd 库,表名加 dwd_ 前缀
  - source-table: ecommerce.\.*
    sink-table: dwd.dwd_<>
    replace-symbol: "<>"
    description: "将 ecommerce.* 路由到 dwd.dwd_*"

上述规则会自动将 ecommerce.orders 路由到 dwd.dwd_orders,ecommerce.users 路由到 dwd.dwd_users。

多库合并:

route:
  - source-table: db0.\.*
    sink-table: merged.<>
    replace-symbol: "<>"
  - source-table: db1.\.*
    sink-table: merged.<>
    replace-symbol: "<>"

8.4 列裁剪与类型转换

Transform 支持:

transform:
  - source-table: ecommerce.orders
    # 列裁剪 + 计算列 + 类型转换
    projection: |
      id,
      CAST(user_id AS STRING) AS user_id_str,
      amount,
      quantity,
      amount * quantity AS total_amount,
      CASE
        WHEN status = '1' THEN 'pending'
        WHEN status = '2' THEN 'paid'
        WHEN status = '3' THEN 'shipped'
        ELSE 'unknown'
      END AS status_desc,
      CURRENT_TIMESTAMP AS sync_time
    filter: amount > 0 AND status IS NOT NULL

3.1 版本支持投影、计算列、常量列;3.2 版本新增 CAST ... AS ... 内置函数,支持 Transform UDF。

8.5 多对一、一对多路由

多对一(分库分表合并):

route:
  # 将多个分表合并到一张表
  - source-table: ecommerce.order_0
    sink-table: ecommerce.all_orders
  - source-table: ecommerce.order_1
    sink-table: ecommerce.all_orders
  - source-table: ecommerce.order_2
    sink-table: ecommerce.all_orders
  # 或使用正则
  - source-table: ecommerce.order_[0-2]
    sink-table: ecommerce.all_orders

当多张分表的 Schema 不完全一致时(如某张表新增了列),Flink CDC 会:

  • 自动将新增列同步到合并表
  • 其他没有该列的表对应字段填充 NULL
  • 类型不一致时尝试推导协变类型(如 FLOAT → DOUBLE,SMALLINT → BIGINT)
  • 无法无损转换时抛出错误停止任务

一对多(广播)(3.2+):

route:
  - source-table: ecommerce.config
    sink-table:
      - ecommerce.config_backup
      - ecommerce.config_archive
      - reporting.config_report

8.6 提交整库同步任务

# 基本提交
bash $FLINK_CDC_HOME/bin/flink-cdc.sh /path/to/mysql-to-doris.yaml

# 指定额外 JAR(如 JDBC 驱动)
bash $FLINK_CDC_HOME/bin/flink-cdc.sh /path/to/pipeline.yaml \
  --jar /path/to/mysql-connector-java-8.0.27.jar

# 从 Savepoint 恢复
bash $FLINK_CDC_HOME/bin/flink-cdc.sh /path/to/pipeline.yaml \
  --from-savepoint /path/to/savepoint

# 提交到 YARN 集群
bash $FLINK_CDC_HOME/bin/flink-cdc.sh /path/to/pipeline.yaml \
  --target yarn-per-job

# 提交到 Kubernetes(3.2+)
bash $FLINK_CDC_HOME/bin/flink-cdc.sh /path/to/pipeline.yaml \
  --target kubernetes-application

8.7 Schema 变更自动同步演示

  1. 初始状态:MySQL orders 表有 id, user_id, amount 三列
  2. 提交 YAML Pipeline 任务,Doris 自动创建对应结构的表
  3. 在 MySQL 中执行 DDL:
-- 新增列
ALTER TABLE orders ADD COLUMN discount DECIMAL(5,2) DEFAULT 0;

-- 修改列类型
ALTER TABLE orders MODIFY COLUMN amount DECIMAL(12,2);

-- 重命名列
ALTER TABLE orders CHANGE COLUMN status order_status VARCHAR(20);
  1. Flink CDC 自动捕获这些 SchemaChangeEvent,在 Doris 中执行对应的 DDL,并继续同步数据,无需重启任务。

注意:Schema Evolution 的效果取决于 Sink 的支持程度。Doris、StarRocks、Paimon 支持较好;Kafka 作为消息队列本身无 Schema 概念,但 Pipeline 会更新内部 Schema 信息。


9. 实时数仓典型架构

9.1 分层架构

基于 Flink CDC 构建实时数仓的典型分层:

在这里插入图片描述

9.2 ODS 层:MySQL CDC → Kafka

ODS 层直接将业务库 Binlog 投递到 Kafka,保留原始变更记录:

-- MySQL CDC Source
CREATE TABLE ods_orders_source (
    id BIGINT,
    user_id BIGINT,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'mysql',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'ecommerce',
    'table-name' = 'orders',
    'server-id' = '5400-5404',
    'scan.startup.mode' = 'initial'
);

-- Kafka Sink
CREATE TABLE ods_orders_kafka (
    id BIGINT,
    user_id BIGINT,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'kafka',
    'topic' = 'ods_orders',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'canal-json'
);

INSERT INTO ods_orders_kafka SELECT * FROM ods_orders_source;

9.3 DWD 层:Kafka → Flink SQL 清洗 → Doris

-- 从 Kafka 读取 ODS 数据
CREATE TABLE dwd_orders_source (
    id BIGINT,
    user_id BIGINT,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'kafka',
    'topic' = 'ods_orders',
    'properties.bootstrap.servers' = 'kafka:9092',
    'properties.group.id' = 'dwd-group',
    'format' = 'canal-json',
    'scan.startup.mode' = 'earliest-offset'
);

-- 维度表:用户信息(CDC 实时更新)
CREATE TABLE dim_users (
    id BIGINT,
    name STRING,
    level STRING,
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'mysql',
    'port' = '3306',
    'username' = 'flinkcdc',
    'password' = 'FlinkCDC@123',
    'database-name' = 'ecommerce',
    'table-name' = 'users',
    'server-id' = '5410-5414'
);

-- DWD 明细宽表写入 Doris
CREATE TABLE dwd_orders_detail (
    order_id BIGINT,
    user_id BIGINT,
    user_name STRING,
    user_level STRING,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = 'doris:8030',
    'table.identifier' = 'dwd.orders_detail',
    'username' = 'root',
    'password' = ''
);

-- 实时维度关联(Lookup Join 或 Regular Join)
INSERT INTO dwd_orders_detail
SELECT
    o.id AS order_id,
    o.user_id,
    u.name AS user_name,
    u.level AS user_level,
    o.product_id,
    o.quantity,
    o.amount,
    o.status,
    o.order_time
FROM dwd_orders_source AS o
LEFT JOIN dim_users FOR SYSTEM_TIME AS OF o.order_time AS u
ON o.user_id = u.id;

9.4 DWS 层:实时聚合

CREATE TABLE dws_sales_summary (
    dt STRING,
    hour_str STRING,
    total_orders BIGINT,
    total_amount DECIMAL(18,2),
    avg_amount DECIMAL(10,2),
    PRIMARY KEY (dt, hour_str) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = 'doris:8030',
    'table.identifier' = 'dws.sales_hourly',
    'username' = 'root',
    'password' = ''
);

INSERT INTO dws_sales_summary
SELECT
    DATE_FORMAT(order_time, 'yyyy-MM-dd') AS dt,
    DATE_FORMAT(order_time, 'HH') AS hour_str,
    COUNT(DISTINCT id) AS total_orders,
    SUM(amount) AS total_amount,
    AVG(amount) AS avg_amount
FROM dwd_orders_detail
GROUP BY
    DATE_FORMAT(order_time, 'yyyy-MM-dd'),
    DATE_FORMAT(order_time, 'HH');

9.5 Lambda vs Kappa 在 CDC 场景

架构描述CDC 场景适用性
Lambda批处理层(T+1 全量修正)+ 速度层(实时增量)双链路早期方案,维护两套代码成本高
Kappa统一流处理,Kafka 保留全量变更日志,重放即可回溯✅ 更适合 CDC,Flink CDC 天然支持全量+增量一体化

Flink CDC 的全量快照 + 增量读取能力使得 Kappa 架构在数据同步场景变得切实可行——无需额外的批处理链路,一套流式代码同时处理历史数据和实时变更。

9.6 CDC 构建实时数仓的优势与挑战

优势:

  • 数据实时性从 T+1 提升到秒级
  • 解耦业务系统与数据系统,业务代码无需改造
  • 一套代码同时处理全量和增量,降低维护成本
  • Flink 强大的 SQL/CEP/窗口能力支撑复杂实时计算

挑战:

  • 业务库的 DDL 变更需要妥善处理(Schema Evolution)
  • 大事务可能导致消费延迟和内存压力
  • 数据质量监控比批处理更难
  • 多表 Join 的状态管理可能成为瓶颈
  • 需要完善的监控和告警体系

10. 性能调优

10.1 全量快照并行度调优

Chunk Key 选择:

  • 优先选择分布均匀的自增主键
  • 复合主键时注意第一列的基数(Cardinality),如果第一列重复值多会导致 Chunk 大小不均
  • 3.2+ 支持指定非主键列作为 Chunk Key,但需确保该列有索引

分片大小调整:

'scan.incremental.snapshot.chunk.size' = '10000'  -- 默认 8096
chunk.size 增大chunk.size 减小
Chunk 数量减少,JM 内存压力小Chunk 数量多,JM 元数据开销大
单 Chunk 数据量大,TM 内存压力大单 Chunk 数据量小,TM 内存压力小
吞吐高,但故障恢复粒度粗吞吐略低,但故障恢复粒度细

经验值:

  • 小表(< 100 万行):默认 8096 即可
  • 大表(> 1 亿行):调大到 20000-50000,减少 Chunk 总数
  • 出现 JM OOM:增大 chunk.size 或增大 jobmanager.memory.heap.size
  • 出现 TM OOM:减小 chunk.size 或增大 taskmanager.memory.framework.heap.size

并行度设置:

SET 'parallelism.default' = '8';

并行度应根据以下因素调整:

  • 数据源的读能力(MySQL 连接数、IO 负载)
  • 下游 Sink 的写入能力
  • 可用的 TaskManager Slot 数量
  • server-id 范围必须大于并行度

10.2 增量阶段 Binlog 消费优化

增量阶段为单并发读取 Binlog,优化方向:

  1. 增大 Debezium 内部队列和批次:
'debezium.max.queue.size' = '162580',   -- 默认 8192
'debezium.max.batch.size' = '40960',    -- 默认 2048
'debezium.poll.interval.ms' = '50'      -- 默认 1000
  1. 心跳保活:
'heartbeat.interval' = '30s'

心跳事件可以:

  • 防止长时间无变更时连接被防火墙/数据库断开
  • 更新低位水印(Low Watermark),避免大事务期间 Checkpoint 超时
  1. 仅解析需要的表(VVR 8.0.7+):
'scan.only.deserialize.captured.tables.changelog.enabled' = 'true'

10.3 服务端 ID 管理与冲突避免

  • 每个运行中的 CDC 任务必须拥有唯一的 server-id
  • 并行读取时 server-id 配置为范围(如 5400-5408),范围大小需大于并行度
  • 不同任务之间不要使用重叠的 server-id 范围
  • 在 MySQL 中查看当前连接的从节点:
SHOW SLAVE HOSTS;

10.4 心跳事件与连接保活

当数据库长时间没有变更时,Binlog 连接可能因空闲超时而断开。配置心跳:

'heartbeat.interval' = '30s'  -- 每 30 秒发送一次心跳

心跳事件的作用:

  • 保持数据库连接活跃
  • 在有未提交大事务时,定期推送位点信息推进 Checkpoint
  • 避免因长时间无数据导致的 Checkpoint 超时

10.5 内存管理

阶段内存消耗点优化建议
快照阶段每个 Chunk 的数据缓存控制 chunk.size,避免单 Chunk 过大
快照阶段JM 中所有 Chunk 元数据大表增大 chunk.size 减少 Chunk 数
增量阶段Binlog 解析队列调整 debezium.max.queue.size
增量阶段大事务事件缓存监控大事务,拆分事务或增大 TM 内存
下游 JoinFlink 状态后端使用 RocksDB 状态后端,增大状态 TTL

Flink 配置建议:

# flink-conf.yaml
jobmanager.memory.process.size: 2048m
taskmanager.memory.process.size: 4096m
taskmanager.memory.framework.heap.size: 512m
state.backend: rocksdb
state.checkpoints.dir: hdfs:///flink/checkpoints
state.savepoints.dir: hdfs:///flink/savepoints
execution.checkpointing.interval: 60000
execution.checkpointing.timeout: 600000

10.6 反压处理

反压信号:在 Flink Web UI 中,Source 或中间算子出现红色 Backpressure 标识。

常见原因与解决方案:

原因表现解决方案
Sink 写入慢Sink 算子反压增大 Sink 并行度、批量写入、优化 Sink 连接
大事务Binlog 阶段突然延迟拆分上游大事务、增大 TM 内存、心跳推进 Checkpoint
数据倾斜某个 Subtask 负载高调整 Chunk Key、重分区、Salting
GC 问题反压与 GC 停顿同时出现增大 JVM 堆、调优 GC 参数、使用 G1GC
下游存储瓶颈Sink 写入延迟高检查目标库负载、增大批量、读写分离

10.7 关键配置参数速查表

参数默认值推荐值说明
scan.incremental.snapshot.enabledtruetrue启用并行无锁快照
scan.incremental.snapshot.chunk.size80968096-50000Chunk 行数
scan.snapshot.fetch.size10242000-5000每次 fetch 行数
server-id随机显式范围必须全局唯一
heartbeat.interval30s30s心跳间隔
connection.pool.size2020-30JDBC 连接池
debezium.max.queue.size8192162580Debezium 内部队列
debezium.max.batch.size204840960每批最大事件数
debezium.poll.interval.ms100050轮询间隔
connect.timeout30s30s连接超时
parallelism.default14-16作业并行度
Checkpoint 间隔—60s生产建议

11. 常见问题与踩坑

11.1 Binlog 格式不对

现象:启动报错 The MySQL server is not configured with binlog_format=ROW 或数据中 UPDATE 事件缺少变更前数据。

原因:binlog_format 设置为 STATEMENT 或 MIXED。

解决:

binlog_format=ROW
binlog_row_image=FULL

动态修改(无需重启,但新连接才生效):

SET GLOBAL binlog_format = 'ROW';
SET GLOBAL binlog_row_image = 'FULL';

11.2 server-id 冲突/重复

现象:报错 A slave with the same server_uuid/server_id as this slave has connected to the master 或数据异常中断。

原因:多个 CDC 任务或 MySQL 从库使用了相同的 server-id。

解决:

  • 为每个 Flink CDC 任务分配不同的 server-id 范围
  • 范围大小大于并行度
  • 记录所有 CDC 任务的 server-id 分配,避免冲突

11.3 全量快照慢/锁表

现象:2.x 版本全量阶段数据库卡住,业务写入阻塞。

原因:旧版快照算法使用 FLUSH TABLES WITH READ LOCK 获取全局读锁。

解决:3.x 默认启用增量快照(scan.incremental.snapshot.enabled=true),无需加锁。如果仍在使用 2.x,升级到 3.x 或确保使用增量快照模式。

11.4 大事务导致延迟/OOM

现象:MySQL 执行一个百万行的大 UPDATE/DELETE 后,Flink CDC 消费延迟飙升甚至 OOM。

原因:Debezium 会将整个事务缓存在内存中,直到事务 COMMIT 才向下游发送。大事务占用大量内存。

解决:

  • 与业务方沟通,将大事务拆分为小批次
  • 增大 TaskManager 内存
  • 配置心跳事件推进 Checkpoint
  • 监控 MySQL 长事务:SELECT * FROM information_schema.innodb_trx

11.5 时区问题

现象:TIMESTAMP / DATETIME 字段时间相差 8 小时。

原因:Flink 集群时区与 MySQL 会话时区不一致。

解决:

'server-time-zone' = 'Asia/Shanghai'

或在 Flink 配置中设置:

env.java.opts: "-Duser.timezone=Asia/Shanghai"

11.6 类型不支持/精度丢失

常见问题:

  • MySQL UNSIGNED BIGINT 超出 Java Long 范围 → 使用 DECIMAL
  • BINARY/VARBINARY → 映射为 BYTES
  • JSON 类型 → 映射为 STRING
  • DECIMAL 精度不一致 → 显式指定 Flink DECIMAL 精度

建议:建表时仔细对照类型映射表,对于特殊类型使用自定义反序列化器处理。

11.7 Schema 变更导致任务失败

现象:上游执行 ALTER TABLE ADD COLUMN 后任务报错。

原因:

  • 2.x 不支持 Schema Evolution
  • 3.x Sink 不支持对应 DDL(如 Kafka Sink 无表结构概念)
  • Schema Evolution 行为模式设置为 exception

解决:

  • 升级到 3.x 并配置 lenient 或 try_evolve 模式
  • 对于 SQL 方式,需要手动在 Flink 表和下游表中添加列后重启任务

11.8 Kafka Connect/Debezium 与 Flink CDC 混用问题

现象:同时使用 Debezium Kafka Connect 和 Flink CDC 读取同一个 MySQL 实例,出现 server-id 冲突或数据重复。

解决:

  • 确保两者使用不同的 server-id
  • 不要将 Flink CDC 与 Kafka Connect Debezium 指向同一个 MySQL server-id
  • 理解两者的定位差异:Flink CDC 是一体化方案,不需要 Kafka Connect 中转

11.9 Exactly-Once 端到端注意事项

端到端 Exactly-Once 需要三个环节同时满足:

  1. Source 端:Flink CDC 基于 Checkpoint 保证 Exactly-Once(✅)
  2. Flink 计算:Checkpoint + State Backend 保证(✅)
  3. Sink 端:需要 Sink 支持事务或幂等写入
Sink 类型Exactly-Once 方式
Kafka两阶段提交(事务)
Doris/StarRocksStream Load Label 幂等
Paimon/Iceberg/HudiSnapshot 两阶段提交
JDBC依赖主键 Upsert 幂等
Redis仅 At-Least-Once

11.10 快照期间 DDL 处理

现象:在全量快照进行期间,上游执行了 DDL,导致快照数据与 Binlog 数据 Schema 不一致。

解决:3.x 的增量快照算法通过 Offset Signal Algorithm 处理快照期间的变更——在 LOW 和 HIGH 水位之间的 Binlog 变更会被 Upsert 到快照结果中,DDL 事件也会被正确处理。但建议尽量避免在全量同步高峰期执行 DDL。

11.11 Checkpoint 失败/状态过大

现象:Checkpoint 频繁超时或失败,State 大小持续增长。

原因与解决:

  • 大事务期间 Checkpoint 超时:配置心跳事件,增大 execution.checkpointing.timeout
  • 状态过大(Join/Aggregation):使用 RocksDB 状态后端,设置 State TTL
  • 快照阶段 Chunk 进度状态过多:增大 chunk.size 减少 Chunk 数量
  • 反压导致 Barrier 对齐慢:解决反压问题,或启用 Unguarded Checkpoint

12. Flink CDC 3.x 新特性深度解析

12.1 增量快照算法演进

Flink CDC 2.0 引入增量快照算法,3.x 持续增强:

  • lock-free:无需 FLUSH TABLES WITH READ LOCK
  • parallel:多 Reader 并行读取不同 Chunk
  • backfill:Chunk 读取完成后,回填 LOW 到 HIGH 水位间的 Binlog 变更
  • checkpoint at chunk granularity:Chunk 级 Checkpoint,故障恢复更高效
  • 3.2 增强:支持任意列作为 Chunk Key(不再局限于主键列)
  • 自动缩容:快照完成后自动减少增量阶段的 Reader 数量,释放资源

12.2 数据源生态扩展

3.x 版本 Pipeline Connector 生态扩展路线:

版本新增 Pipeline Connector
3.0MySQL Source, Doris Sink, StarRocks Sink
3.1Kafka Sink, Paimon Sink
3.2Elasticsearch Sink
3.5Fluss Sink, PostgreSQL Source
3.6Oracle Source, Hudi Sink

SQL Source Connector 已覆盖 10+ 种数据库。

12.3 Pipeline Connector 与整库同步 YAML

YAML Pipeline 是 3.x 的核心创新:

  • 声明式配置:无需写 Java/Scala 代码
  • 自动建表:下游表自动创建,无需手动 DDL
  • 端到端集成:Source → Transform → Route → Sink 全链路
  • CLI 提交:flink-cdc.sh 一行命令
  • 多部署模式:Standalone、YARN、Kubernetes(3.2+)

12.4 Schema Evolution 演进

版本Schema Evolution 能力
3.0基础 AddColumn / DropColumn 支持
3.1Transform 与 Schema Evolution 协同
3.2五种行为模式(Lenient / TryEvolve / Evolve / Ignore / Exception)+ 事件类型级控制
3.5故障恢复时重新下发 Schema 信息;大小写敏感处理
3.6PostgreSQL Schema Evolution 支持;AlterTableComment / AlterColumnComment 事件;投影更新时 Schema 差异检测

12.5 动态新增表同步

3.x 支持在作业运行过程中动态添加监控表:

  • 修改 YAML 配置中的 tables 规则
  • 从 Savepoint 恢复任务
  • 新添加的表会自动触发全量快照并接入增量 Binlog 流
  • 不影响已有表的同步进度

12.6 Flink 版本兼容性

Flink CDC 版本支持的 Flink 版本JDK 要求
3.0.x1.14 - 1.18JDK 8+
3.1.x1.14 - 1.18JDK 8+
3.2.x1.17 - 1.20JDK 8+
3.5.x1.18 - 1.20JDK 8+
3.6.x1.20.x / 2.xJDK 11(强制)

12.7 未来路线图

根据 Flink CDC 社区的公开讨论,未来方向包括:

  • 更多 Pipeline Connector(如更多 Sink 支持)
  • Transform 能力增强(更多内置函数、复杂类型支持)
  • 3.6 已引入 VARIANT 类型和 PARSE_JSON 函数,增强半结构化数据处理
  • 正则路由(3.6+)和首个匹配路由规则
  • 更完善的 Schema Evolution(如 PostgreSQL 已在 3.6 支持)
  • 性能持续优化(Binlog 读取、Chunk 切分、并行解析)

13. 端到端实战项目

13.1 项目背景

某电商平台业务库(MySQL)包含订单、用户、商品三张核心表,需要实时同步到 Doris 数仓,支撑实时大屏和经营分析。

13.2 技术栈

MySQL → Flink CDC (YAML Pipeline) → Kafka (ODS)
                                      ↓
                               Flink SQL (DWD/DWS)
                                      ↓
                               Doris (ADS)

13.3 表结构准备

-- MySQL 建表
CREATE DATABASE ecommerce;
USE ecommerce;

CREATE TABLE users (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    name VARCHAR(100) NOT NULL,
    email VARCHAR(200),
    level VARCHAR(20) DEFAULT 'bronze',
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

CREATE TABLE products (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    name VARCHAR(200) NOT NULL,
    category VARCHAR(50),
    price DECIMAL(10,2),
    stock INT DEFAULT 0,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

CREATE TABLE orders (
    id BIGINT AUTO_INCREMENT PRIMARY KEY,
    user_id BIGINT NOT NULL,
    product_id BIGINT NOT NULL,
    quantity INT NOT NULL,
    amount DECIMAL(10,2) NOT NULL,
    status VARCHAR(20) DEFAULT 'pending',
    order_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    INDEX idx_user (user_id),
    INDEX idx_product (product_id)
);

-- 插入测试数据
INSERT INTO users (name, email, level) VALUES
('Alice', 'alice@example.com', 'gold'),
('Bob', 'bob@example.com', 'silver'),
('Carol', 'carol@example.com', 'bronze');

INSERT INTO products (name, category, price, stock) VALUES
('iPhone 15', 'electronics', 6999.00, 100),
('AirPods Pro', 'electronics', 1899.00, 200),
('Python Book', 'books', 89.00, 500);

INSERT INTO orders (user_id, product_id, quantity, amount, status) VALUES
(1, 1, 1, 6999.00, 'paid'),
(2, 2, 2, 3798.00, 'shipped'),
(1, 3, 3, 267.00, 'pending');

13.4 整库同步 MySQL 到 Kafka(ODS)

创建 mysql-to-kafka.yaml:

pipeline:
  name: Ecommerce-ODS-Sync
  parallelism: 3
  schema.change.behavior: lenient

source:
  type: mysql
  hostname: 127.0.0.1
  port: 3306
  username: flinkcdc
  password: FlinkCDC@123
  tables: ecommerce.\.*
  server-id: 5400-5406
  server-time-zone: Asia/Shanghai
  scan.startup.mode: initial

sink:
  type: kafka
  properties.bootstrap.servers: 127.0.0.1:9092
  topic: ods_ecommerce
  key.format: json
  value.format: debezium-json
  partition.strategy: hash-by-key

route:
  - source-table: ecommerce.\.*
    sink-table: ods_ecommerce

提交:

bash flink-cdc.sh mysql-to-kafka.yaml

13.5 Flink SQL 清洗写入 Doris DWD

-- ============================================
-- DWD:订单明细宽表
-- ============================================
CREATE TABLE kafka_ods_orders (
    id BIGINT,
    user_id BIGINT,
    product_id BIGINT,
    quantity INT,
    amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'kafka',
    'topic' = 'ods_ecommerce',
    'properties.bootstrap.servers' = '127.0.0.1:9092',
    'properties.group.id' = 'dwd-orders',
    'format' = 'debezium-json',
    'scan.startup.mode' = 'earliest-offset'
);

CREATE TABLE kafka_dim_users (
    id BIGINT,
    name STRING,
    level STRING,
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'kafka',
    'topic' = 'ods_ecommerce',
    'properties.bootstrap.servers' = '127.0.0.1:9092',
    'properties.group.id' = 'dwd-users',
    'format' = 'debezium-json',
    'scan.startup.mode' = 'earliest-offset'
);

CREATE TABLE doris_dwd_order_detail (
    order_id BIGINT,
    user_id BIGINT,
    user_name STRING,
    user_level STRING,
    product_id BIGINT,
    product_name STRING,
    product_category STRING,
    quantity INT,
    order_amount DECIMAL(10,2),
    status STRING,
    order_time TIMESTAMP(3),
    PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = '127.0.0.1:8030',
    'table.identifier' = 'dwd.order_detail',
    'username' = 'root',
    'password' = '',
    'sink.label-prefix' = 'dwd_order'
);

INSERT INTO doris_dwd_order_detail
SELECT
    o.id,
    o.user_id,
    u.name,
    u.level,
    o.product_id,
    p.name,
    p.category,
    o.quantity,
    o.amount,
    o.status,
    o.order_time
FROM kafka_ods_orders o
LEFT JOIN kafka_dim_users FOR SYSTEM_TIME AS OF o.order_time u
    ON o.user_id = u.id
LEFT JOIN kafka_dim_products FOR SYSTEM_TIME AS OF o.order_time p
    ON o.product_id = p.id;

13.6 实时聚合写入 Doris DWS

CREATE TABLE doris_dws_hourly_sales (
    stat_hour STRING,
    category STRING,
    order_count BIGINT,
    total_quantity BIGINT,
    total_amount DECIMAL(18,2),
    PRIMARY KEY (stat_hour, category) NOT ENFORCED
) WITH (
    'connector' = 'doris',
    'fenodes' = '127.0.0.1:8030',
    'table.identifier' = 'dws.hourly_sales',
    'username' = 'root',
    'password' = '',
    'sink.label-prefix' = 'dws_sales'
);

INSERT INTO doris_dws_hourly_sales
SELECT
    DATE_FORMAT(order_time, 'yyyy-MM-dd-HH') AS stat_hour,
    product_category AS category,
    COUNT(DISTINCT order_id) AS order_count,
    SUM(quantity) AS total_quantity,
    SUM(order_amount) AS total_amount
FROM doris_dwd_order_detail
GROUP BY
    DATE_FORMAT(order_time, 'yyyy-MM-dd-HH'),
    product_category;

13.7 数据验证与监控

数据一致性验证:

-- MySQL 端订单总数
SELECT COUNT(*) FROM ecommerce.orders;

-- Doris 端订单总数
SELECT COUNT(*) FROM dwd.order_detail;

-- 对比金额
SELECT SUM(amount) FROM ecommerce.orders;
SELECT SUM(order_amount) FROM dwd.order_detail;

监控指标:

  • Flink Web UI:Checkpoint 成功率、反压、消费延迟
  • 自定义指标:Source 端 recordsConsumed、Sink 端 recordsWritten
  • Doris:Stream Load 成功率、延迟
  • 消费延迟:MysqlSourceEventFetchDelay 或通过 Binlog 位点与当前时间差计算

13.8 完整 YAML 一键整库同步(更简洁的替代方案)

实际上,3.x 的 YAML Pipeline 可以直接完成 MySQL → Doris 的整库同步,无需上面的中间 Kafka 步骤:

pipeline:
  name: Ecommerce-Full-Sync
  parallelism: 4
  schema.change.behavior: lenient

source:
  type: mysql
  hostname: 127.0.0.1
  port: 3306
  username: flinkcdc
  password: FlinkCDC@123
  tables: ecommerce.\.*
  server-id: 5400-5408
  server-time-zone: Asia/Shanghai

sink:
  type: doris
  fenodes: 127.0.0.1:8030
  username: root
  password: ""
  table.create.properties.light_schema_change: "true"
  table.create.properties.replication_num: "1"

transform:
  - source-table: ecommerce.orders
    projection: id, user_id, product_id, quantity, amount, status, order_time,
                amount * quantity AS total_amount
    filter: status IS NOT NULL

route:
  - source-table: ecommerce.\.*
    sink-table: dwd.<>
    replace-symbol: "<>"
bash flink-cdc.sh ecommerce-full-sync.yaml

一行命令完成:全量快照 → 增量同步 → Schema Evolution → 自动建表 → 数据变换 → 表名路由。这就是 3.x YAML Pipeline 的威力。


14. 学习路线与资源

14.1 学习阶段

阶段一:入门(1-2 周)
├── 理解 CDC 概念和 Binlog 原理
├── 搭建本地环境(Docker Compose)
├── 运行第一个 Flink SQL CDC Demo
└── 掌握 MySQL CDC Source DDL 参数

阶段二:进阶(2-4 周)
├── DataStream API 自定义反序列化
├── 配置多种 Sink(Kafka/Doris/Paimon)
├── 理解增量快照算法原理
├── Checkpoint 和 Exactly-Once 机制
└── 性能调优基础

阶段三:实战(4-8 周)
├── YAML Pipeline 整库同步
├── Transform 和 Route 配置
├── Schema Evolution 实践
├── 实时数仓分层建设
└── 生产问题排查

阶段四:源码(持续)
├── Flink CDC 源码编译与调试
├── 增量快照算法源码阅读
├── 自定义 Connector 开发
└── 社区贡献

14.2 官方资源

资源链接
官方文档https://nightlies.apache.org/flink/flink-cdc-docs-stable/
GitHub 仓库https://github.com/apache/flink-cdc
Release 下载https://github.com/apache/flink-cdc/releases
Flink 官网https://flink.apache.org/
Debezium 文档https://debezium.io/documentation/
DBLog 论文https://arxiv.org/abs/2610.01724
Issue 跟踪https://issues.apache.org/jira/projects/FLINK/summary

14.3 面试高频考点

1. Flink CDC 的全量和增量阶段是如何切换的?

所有 Chunk 快照读取完成后,等待一次完整 Checkpoint 成功(确保快照数据全部被下游消费),然后从快照结束时记录的 Binlog 高水位位点开始增量读取。

2. 增量快照算法为什么不需要加锁?

它使用 Offset Signal Algorithm:先记录 LOW 位点,读取 Chunk 数据,再记录 HIGH 位点,然后读取 LOW-HIGH 之间的 Binlog 变更并 Upsert 到快照结果中。整个过程通过 Binlog 回填保证一致性,无需全局读锁。

3. server-id 的作用和配置注意事项?

server-id 是 MySQL 复制协议中从节点的唯一标识。每个 CDC 任务必须使用唯一的 server-id;并行读取时需配置为范围,范围大小要大于并行度。

4. Flink CDC 如何保证 Exactly-Once?

Source 端基于 Chunk 级 Checkpoint 和 Binlog 位点 State 保证不丢不重;Flink 引擎通过 Checkpoint 和两阶段提交保证计算一致性;Sink 端需要支持事务或幂等写入。

5. 增量阶段为什么是单并发?

Binlog 是有序的事件日志,必须按顺序消费才能保证数据一致性。多并发消费会导致事件乱序,产生错误的最终状态。

6. 2.x 和 3.x 的主要区别?

3.x 引入了 YAML Pipeline 整库同步、Schema Evolution、Transform/Route、Pipeline Connector 生态、自动缩容、动态加表等端到端数据集成能力,定位从 Connector 集合升级为数据集成框架。

7. 如何处理 Schema 变更?

3.x 使用 Schema Operator + Schema Coordinator,在 DDL 事件到达时先 Flush 下游数据,再通过 MetadataApplier 应用 DDL。支持 lenient/try_evolve/evolve/ignore/exception 五种行为模式。

8. Chunk Key 如何选择?

优先选择分布均匀的自增主键。复合主键注意第一列的基数。3.2+ 支持非主键列但需有索引。无主键表必须指定 chunk.key-column。

9. 大事务导致 OOM 怎么处理?

拆分上游大事务;增大 TM 内存;配置心跳推进 Checkpoint;监控长事务。

10. Flink CDC vs Debezium 的区别?

Flink CDC 底层用 Debezium 解析 Binlog,但增强了并行快照、Checkpoint 集成、SQL/YAML API、整库同步、Schema Evolution、端到端 ETL 能力。Debezium 依赖 Kafka Connect,单 Task 限制明显。

11. 全量快照慢如何优化?

增大并行度和 chunk.size;优化 Chunk Key 选择;增大 fetch.size 和连接池;检查 MySQL 负载和索引。

12. Canal-JSON 和 Debezium-JSON 格式区别?

Canal-JSON 是阿里巴巴 Canal 的格式,使用 data/old/type 字段,国内生态友好;Debezium-JSON 使用 before/after/op/source 字段,是 Debezium 标准格式,国际通用。

13. 动态加表的原理?

修改监控表配置后从 Savepoint 恢复,Enumerator 发现新表后为其切分 Chunk 并执行快照,同时在 Binlog 流中过滤该表的变更,快照完成后接入增量流。

14. 分库分表合并时 Schema 不一致怎么办?

Flink CDC 自动合并 Schema:新增列同步到下游,缺失列填 NULL;类型不一致时尝试协变类型(无损转换);无法无损处理时报错停止。

15. Flink CDC 支持哪些数据源和目标端?

Source 支持 MySQL、PostgreSQL、Oracle、SQL Server、MongoDB、DB2、OceanBase、TiDB、Vitess;Pipeline Sink 支持 Doris、StarRocks、Kafka、Paimon、Elasticsearch、Iceberg、Hudi、Fluss。

14.4 实战练习建议

  1. 基础练习:用 Docker Compose 搭建 MySQL + Flink,用 SQL Client 实现单表 CDC 到 MySQL
  2. 进阶练习:配置 MySQL CDC → Kafka → Flink SQL → Doris 的完整链路
  3. YAML 练习:用 YAML Pipeline 实现 MySQL 整库同步到 Doris,体验自动建表
  4. Schema Evolution 练习:在同步过程中对 MySQL 表执行 ADD COLUMN / MODIFY COLUMN,观察下游自动变更
  5. 分库分表练习:创建 order_0/order_1/order_2,配置 Route 合并到一张表
  6. 性能练习:构造千万级数据表,调整并行度和 chunk.size,观察吞吐量变化
  7. 容错练习:同步过程中手动 kill TaskManager,观察 Checkpoint 恢复和数据一致性

15. 总结

Flink CDC 已经从一个简单的 MySQL Binlog 读取工具,演进为功能完备的分布式实时数据集成框架。它的核心价值可以概括为:

  1. 全量增量一体化:无需在 DataX(全量)和 Canal(增量)之间切换,一套代码搞定
  2. 无锁并行快照:对生产业务库零侵入,水平扩展快照吞吐量
  3. SQL + YAML 双 API:SQL 适合单表精细 ETL,YAML 适合整库快速同步
  4. Schema Evolution:上游 DDL 自动同步,不再因加列改类型而重启任务
  5. 丰富生态:10+ Source,8+ Sink,覆盖主流数据库和数据湖/数仓
  6. 生产可靠:Checkpoint 容错、Exactly-Once、动态扩缩容

对于数据开发工程师来说,掌握 Flink CDC 不仅是学会一个工具,更是理解实时数据集成范式转变的关键——从 T+1 批量到秒级流式,从代码硬编码到声明式配置,从手动 Schema 管理到自动演化。

希望这篇博客能成为你学习 Flink CDC 的起点。建议按照第 14 章的学习路线,先在本地跑通 Demo,再逐步深入原理和生产实践。遇到问题时,官方文档和 GitHub Issue 是最好的朋友。

最后的话:技术在快速演进,Flink CDC 3.2 不是终点,3.5、3.6 已经带来更多连接器和能力。保持学习,关注社区,动手实践——这是每个优秀数据工程师的成长之路。


本文基于 Flink CDC 3.2.1 版本撰写,部分内容参考 3.5/3.6 版本特性。由于开源项目迭代迅速,具体参数和行为请以对应版本的官方文档为准。


更多推荐