传统的数仓 ETL 通常依赖定时调度任务,比如每天凌晨用 Sqoop 或者 DataX 全量抽取一次业务库数据。这种方式实现简单,但延迟以小时甚至天为单位,一旦业务方要求看分钟级的实时大盘,批处理方案就完全撑不住了。Flink CDC 的出现改变了这个局面,它能够直接订阅 MySQL 的 binlog,把每一行数据的变更当作流事件来处理,再结合 Flink 强大的流计算能力,一条从业务库到数仓的秒级同步链路就搭建起来了。本文将以 Spring Boot 为载体,完整演示如何整合 Flink CDC,把 MySQL 中的业务数据实时同步到分析型存储中。

一、先弄懂 Flink CDC 的工作原理
CDC 的全称是 Change Data Capture,即变更数据捕获。它的核心思想是:数据库每一次 insert、update、delete 操作都会被记录在 binlog 里,只要能消费 binlog,就能拿到数据的全量变更历史。Flink CDC 连接器(flink-connector-mysql-cdc)内部基于 Debezium 实现,Debezium 负责解析 binlog 并转成统一的事件格式,Flink 负责把这些事件当作无界流来计算。
Flink CDC 最大的亮点是无锁快照机制。早期版本做全量初始化时需要对表加全局锁,会阻塞线上写入。新版连接器借鉴了 Netflix DBLog 论文的方案,通过分块读取加高低水位位的算法,在不加锁的情况下保证全量数据和增量 binlog 的无缝衔接,既不锁表,也不会丢数据或者数据重复。理解这一点很重要,因为直接决定了这条链路能否用在生产库上。
整个链路的数据流转路径是:MySQL binlog 被捕获后,先执行全量快照读取存量数据,再切换到增量模式持续消费日志;流经过 Flink 的转换算子完成清洗、维表关联、打宽等操作,最后 sink 到目标端,比如 Doris、StarRocks、Kafka 或者另一个 MySQL。这就是一条完整的实时 ETL 链路。
二、环境准备与 MySQL 端配置
整合之前要先把基础环境准备好。Flink CDC 对 MySQL 有几个硬性要求:binlog 格式必须是 ROW,binlog_row_image 必须是 FULL,同步账号需要授予 SELECT 和 REPLICATION SLAVE、REPLICATION CLIENT 权限。先检查数据库配置:
-- 检查 binlog 配置 SHOW VARIABLES LIKE 'log_bin'; -- 需要为 ON SHOW VARIABLES LIKE 'binlog_format'; -- 需要为 ROW SHOW VARIABLES LIKE 'binlog_row_image'; -- 需要为 FULL -- 创建专用同步账号 CREATE USER 'flink_user'@'%' IDENTIFIED BY 'your_password'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flink_user'@'%'; FLUSH PRIVILEGES;
如果 log_bin 显示为 OFF,需要修改 my.cnf,在 mysqld 段落下加上 log-bin=mysql-bin、binlog_format=ROW 和 server_id 配置后重启。另外要特别注意,同步账号的密码策略和 host 白名单,很多连不上的问题都出在这里。
Flink 侧需要部署一个独立集群(推荐 standalone 或 on YARN 模式),并把 mysql-cdc 连接器的 jar 包放到 Flink 的 lib 目录下。连接器版本要和 Flink 版本严格匹配,比如 Flink 1.17 搭配 flink-sql-connector-mysql-cdc 2.4.x,版本错配会直接抛 ClassNotFound 异常,这是新手最常见的坑之一。
三、Spring Boot 侧的两种整合方式
Spring Boot 和 Flink CDC 的整合,常见有两种思路。第一种是进程外提交方式:Spring Boot 只负责管理和触发任务,真正的 Flink 作业通过 SQL Client 或者 REST API 提交到 Flink 集群运行,计算逻辑与业务系统完全解耦。第二种是嵌入式执行方式:把 Flink 作为普通依赖引入 Spring Boot 项目,用 DataStream API 在本地 MiniCluster 上跑。两种方式各有适用场景,下面分别展开。
方式一:DataStream API 嵌入式开发
先引入依赖,注意 flink 相关依赖要设置 scope 为 provided,避免和 Spring Boot 自带的依赖冲突:
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>2.4.2</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.2</version>
</dependency>
<!-- Spring Boot 侧还需排除自带的 flink 相关传递依赖,防止类冲突 -->
接着在 Service 层编写同步逻辑。下面的例子把 MySQL 的订单表变更实时打印并写出到 Kafka,实际项目中可以把 Kafka 换成 Doris 的 connector:
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class OrderCdcSync {
public void start() throws Exception {
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("127.0.0.1")
.port(3306)
.databaseList("order_db")
.tableList("order_db.t_order") // 必须带上库名前缀
.username("flink_user")
.password("your_password")
.deserializer(new JsonDebeziumDeserializationSchema())
.startupOptions(StartupOptions.initial()) // 先全量后增量
.build();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60000); // 开启 checkpoint,保证 exactly-once
env.fromSource(source, WatermarkStrategy.noWatermarks(), "mysql-cdc-source")
.map(this::transform) // 这里做清洗、转换、打宽
.sinkTo(buildKafkaSink());
env.execute("order-cdc-sync");
}
}
有几个关键点要说明。第一,startupOptions 默认是 initial,即先做全量快照再追增量,生产上首次上线基本都用这个模式。第二,checkpoint 必须开启,Flink CDC 依赖 checkpoint 机制来持久化 binlog 位点,如果不开 checkpoint,作业重启后会从头消费,代价非常大。第三,嵌入式方式只适合小规模数据量的场景,Spring Boot 进程里跑 Flink 会互相争抢资源,且无法享受集群的高可用能力。
方式二:SQL 方式提交到集群
生产环境更推荐用 Flink SQL 来写同步任务,Spring Boot 通过 REST API 把作业提交到 Flink 集群。SQL 写法极大降低了开发成本,比如把 MySQL 数据同步到 Doris 只需要创建源表和目标表再插入即可:
-- 创建 CDC 源表
CREATE TABLE orders_source (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10, 2),
status STRING,
create_time TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '127.0.0.1',
'port' = '3306',
'username' = 'flink_user',
'password' = 'your_password',
'database-name' = 'order_db',
'table-name' = 't_order'
);
-- 写入 Doris 目标表
INSERT INTO doris_orders
SELECT order_id, user_id, amount, status, create_time FROM orders_source;
Spring Boot 侧的职责则变成任务的生命周期管理:把上述 SQL 保存在配置中心,通过 Flink REST 接口的 /jars/:jobid/run 触发提交,通过 /jobs/:jobid 查询状态,作业失败时自动告警和重启。这种架构下,同步链路和业务系统完全解耦,Flink 集群可以独立扩缩容,是最稳妥的生产方案。
四、生产环境的调优点与常见坑
链路跑通只是第一步,真正上生产还需要关注几个方面。首先是并行度与分块读取。全量阶段 CDC 会按照主键把表切成多个 chunk 并行读取,可以通过 scan.incremental.snapshot.chunk-size 控制每块大小,大表建议设置成 80640 左右,避免单个 chunk 过大导致读取超时。
其次是位点与容错。一定要配置 checkpoint 存储(HDFS 或 S3)并设置合理的间隔,同时把重启策略设置为 failure-rate,避免数据库瞬时抖动导致作业永久退出。如果源库存在批量删除或大事务,binlog 事件会瞬间暴增,要适当调大 taskmanager 的内存,否则容易 OOM。
最后是一些容易踩的坑:表必须有主键,否则增量阶段无法定位变更行;同步的表结构变更(DDL)在旧版本连接器中不支持自动同步,加字段前要先停作业;时区问题也很常见,建议在连接器参数里显式配置 server-time-zone 为 Asia/Shanghai,否则时间字段会差 8 小时。把这些细节处理到位,一条稳定支撑每天数十亿变更事件的实时 ETL 链路就真正落地了。
Spring BootFlink CDC实时数据仓库修改时间:2026-09-06 21:18:45