第1章 Flink CDC 概述
1.1 什么是 CDC(变更数据捕获)
| 概念名称 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
| CDC(Change Data Capture) | 一种用于捕获数据库中数据变更(插入、更新、删除)的技术,通常通过监听数据库的事务日志(如 MySQL 的 binlog)实现。 | 实现数据的实时同步、数据复制、缓存更新、事件驱动架构等。 | 需要数据库开启日志功能(如 binlog),并确保日志格式为 ROW 模式。 |
| Binlog(Binary Log) | MySQL 中记录所有数据变更操作的日志文件,是 CDC 实现的基础。 | 提供数据变更的原始记录,供下游系统消费。 | 必须配置 binlog_format=ROW,否则无法捕获字段级变更。 |
| Debezium | 开源分布式 CDC 平台,基于 Kafka Connect 构建,Flink CDC 内部集成其核心引擎。 | 提供统一的变更事件格式和连接器支持。 | Flink CDC 封装了 Debezium,无需独立部署 Kafka Connect。 |
| Snapshot(快照) | 在首次启动时,Flink CDC 会先对表进行全量快照读取,再切换到 binlog 流式读取。 | 确保历史数据不丢失,实现全量 + 增量一体化同步。 | 大表快照可能影响源库性能,建议在低峰期执行。 |
1.2 Flink CDC 的核心优势与应用场景
| 概念名称 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
| 全量 + 增量一体化 | 支持首次全量读取表数据,自动切换到增量 binlog 模式,无需手动切换。 | 简化数据同步流程,避免数据断层。 | 需保证 binlog 保留时间足够长,覆盖快照完成时间。 |
| 无锁读取(Lock-free Snapshot) | 使用 RR 隔离级别和 SELECT ... FOR SHARE 或 READ UNCOMMITTED 实现快照,避免锁表。 | 减少对生产库的影响,提升并发能力。 | 需数据库支持 MVCC(如 InnoDB)。 |
| 端到端 Exactly-Once 语义 | 结合 Flink Checkpoint 机制,确保每条变更仅被处理一次。 | 保障数据一致性,适用于金融、订单等关键业务。 | Sink 端需支持幂等写入或事务。 |
| Schema Evolution 支持 | 支持表结构变更(如加列)的自动感知(部分支持)。 | 适应业务迭代,减少运维成本。 | 当前对 DDL 的完整支持仍在演进中。 |
| 应用场景:实时数仓 | 将业务库数据实时同步到数仓(如 Doris、ClickHouse、Iceberg)。 | 构建低延迟的数据分析平台。 | 需注意目标系统写入性能。 |
| 应用场景:缓存更新 | 捕获数据库变更后,自动更新 Redis 缓存。 | 避免缓存与数据库不一致。 | 需设计合理的缓存失效策略。 |
| 应用场景:微服务解耦 | 将数据库变更作为事件发布,驱动其他服务响应。 | 实现事件驱动架构(EDA)。 | 需结合消息中间件(如 Kafka)进行解耦。 |
1.3 Flink CDC 与传统数据同步方案对比
| 对比项 | Flink CDC | 传统方案(如 DataX、Sqoop) | 说明 | 注意事项 |
|---|---|---|---|---|
| 同步模式 | 实时流式同步(秒级延迟) | 批处理同步(分钟/小时级延迟) | Flink CDC 基于流处理,延迟更低。 | 传统方案适合离线批处理场景。 |
| 架构复杂度 | 简单:Flink Job 直接读取数据库日志 | 复杂:常需 Kafka、Kafka Connect 中转 | Flink CDC 内嵌 Debezium,无需额外组件。 | 减少运维成本。 |
| 数据一致性 | 支持端到端 Exactly-Once | 通常为 At-Least-Once,易重复 | Flink Checkpoint 保障一致性。 | Sink 需支持事务或幂等。 |
| 全量增量一体化 | 支持自动切换 | 需分别配置全量和增量任务 | 简化任务管理。 | 避免数据重复或遗漏。 |
| 资源占用 | 持续运行,占用一定资源 | 按需运行,资源占用低 | 适合长期运行的实时任务。 | 需合理配置资源。 |
| 扩展性 | 支持并行读取、多表同步 | 扩展性差,通常单线程 | 可水平扩展应对大数据量。 | 需合理设置分片参数。 |
1.4 Flink CDC 支持的数据源概览(MySQL、PostgreSQL、Oracle 等)
| 数据源 | 支持状态 | 所需依赖 | 关键特性 | 注意事项 |
|---|---|---|---|---|
| MySQL | 官方稳定支持 | flink-connector-mysql-cdc | 支持分库分表、GTID、SSL、读取视图等 | 需开启 binlog,格式为 ROW,server-id 唯一 |
| PostgreSQL | 官方稳定支持 | flink-connector-postgres-cdc | 基于逻辑复制槽(replication slot),支持 toast 超长字段 | 需启用 wal_level=logical,用户有 replication 权限 |
| Oracle | 官方支持(社区版) | flink-connector-oracle-cdb-cdc | 基于 LogMiner 或 XStream,支持多租户 | 配置复杂,需归档日志(archivelog)模式 |
| SQL Server | 社区支持 | flink-connector-sqlserver-cdc | 基于 CDC 功能或变更跟踪 | 需启用数据库级和表级 CDC |
| MongoDB | 社区支持 | flink-connector-mongodb-cdc | 基于 oplog 监听变更 | 需 replica set 或 sharded cluster |
| TiDB | 兼容 MySQL 协议 | 使用 MySQL 连接器 | 可通过 MySQL 模式接入 | 需确认 binlog 开启(如 Pump/Drainer) |
| Cassandra | 实验性支持 | flink-connector-cassandra-cdc | 基于 commitlog | 社区活跃度较低,功能有限 |
注: 所有连接器均可通过 Maven 引入,具体坐标见第2章。
第2章 环境准备与快速上手
2.1 开发环境搭建(Java/Scala + Maven)
| 工具 | 版本要求 | 安装方式 | 用途 | 注意事项 |
|---|---|---|---|---|
| Java | JDK 8 或 JDK 11 | 官网下载并配置 JAVA_HOME | Flink 运行基础环境 | 推荐使用 OpenJDK,避免商业授权问题 |
| Maven | 3.5+ | 官网下载并配置 PATH | 项目依赖管理与构建 | 配置国内镜像(如阿里云)提升下载速度 |
| IDE | IntelliJ IDEA / Eclipse | 官网下载 | 代码编写与调试 | 推荐 IDEA,对 Flink 有良好支持 |
| Flink | 1.13+(推荐 1.17+) | 官网下载或通过 Maven 依赖 | 流处理引擎 | Flink CDC 3.x 要求 Flink 1.17+ |
| 构建命令 | mvn clean package | 命令行执行 | 打包项目为 jar | 确保 pom.xml 正确配置依赖 |
2.2 Flink 运行环境配置(Local/Standalone/YARN)
| 运行模式 | 配置方式 | 适用场景 | 启动命令 | 注意事项 |
|---|---|---|---|---|
| Local 模式 | 无需配置,直接运行 main 方法 | 本地开发测试 | 直接在 IDE 中运行 | 仅用于学习和调试 |
| Standalone 集群 | 解压 Flink 包,修改 conf/flink-conf.yaml | 独立集群部署 | ./bin/start-cluster.sh | 需配置 jobmanager.rpc.address |
| YARN 模式 | 确保 Hadoop 环境可用 | 企业级资源调度 | yarn-session.sh + flink run | 需上传 jar 到 HDFS 或本地 |
| Application 模式 | 使用 flink run-application | 生产推荐,资源隔离 | flink run-application -t yarn-application | Jar 包需包含依赖(fat jar) |
| Session 模式 | 先启动集群,再提交任务 | 多任务共享集群 | flink run -d job.jar | 任务间可能资源竞争 |
2.3 添加 Flink CDC 依赖(Maven 坐标)
| 依赖名称 | 用途 | 注意事项 |
|---|---|---|
| Flink Java API | 提供 DataStream API | 版本需与 Flink 集群一致 |
| Flink Streaming | 流处理核心库 | 通常与 flink-java 一起引入 |
| Flink CDC MySQL | 支持 MySQL 数据源 | Ververica 为官方维护团队 |
| Flink CDC PostgreSQL | 支持 PostgreSQL 数据源 | 需额外配置 replication slot |
| Kafka 连接器(可选) | 将 CDC 数据写入 Kafka | 需引入 kafka-clients |
| Scala 依赖(如用 Scala) | Scala API 支持 | 注意 Scala 版本匹配(2.11/2.12) |
Maven 坐标示例:
<!-- Flink Java API -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.17.0</version>
</dependency>
<!-- Flink Streaming -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.0</version>
</dependency>
<!-- Flink CDC MySQL -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>3.0.1</version>
</dependency>
<!-- Flink CDC PostgreSQL -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-postgres-cdc</artifactId>
<version>3.0.1</version>
</dependency>
<!-- Kafka 连接器(可选) -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.17.0</version>
</dependency>
<!-- Scala 依赖(如用 Scala) -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-scala_2.12</artifactId>
<version>1.17.0</version>
</dependency>
提示: 使用
<scope>provided</scope>避免将 Flink 核心包打入 jar(生产部署时由集群提供)。
2.4 第一个 Flink CDC 程序:MySQL 到控制台
| 方法/配置 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
MySqlSource.builder() | MySqlSource.<String>builder() | 创建 MySQL Source 构建器 | MySqlSource.builder() | 泛型指定输出数据类型(如 String、RowData) |
hostname() | .hostname("localhost") | 设置数据库地址 | .hostname("192.168.0.1") | 可为 IP 或域名 |
port() | .port(3306) | 设置数据库端口 | .port(3306) | 默认 3306 |
databaseName() | .databaseName("mydb") | 指定监听的数据库 | .databaseName("inventory") | 支持正则表达式(如 "test_db.*") |
tableName() | .tableName("users") | 指定监听的表 | .tableName("user_info") | 支持正则(如 "user_.*") |
username() / password() | .username("user").password("pass") | 认证信息 | .username("flink").password("cdc123") | 建议使用只读账号 |
startupMode() | .startupMode(StartupMode.INITIAL) | 启动模式 | .startupMode(StartupMode.LATEST_OFFSET) | INITIAL=全量+增量,LATEST_OFFSET=仅增量 |
debeziumConfig() | .debeziumConfig("snapshot.locking.mode", "none") | 传递 Debezium 配置 | .debeziumConfig("binlog.buffer.size", "16384") | 高级调优参数 |
build() | .build() | 构建 Source 实例 | MySqlSource source = builder.build(); | 返回 SourceFunction |
| 执行 Job | env.addSource(source) | 添加 Source 并启动 | 完整示例见下方 | Source 并行度通常为 1(单分片) |
完整代码示例(Java):
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
MySqlSource mySqlSource = MySqlSource.builder()
.hostname("localhost")
.port(3306)
.databaseName("test_db")
.tableName("user_info")
.username("flink")
.password("cdc123")
.startupMode(StartupMode.INITIAL)
.deserializer(ChangelogJsonDeserializationSchema.INSTANCE)
.build();
env.addSource(mySqlSource, "MySQL Source")
.setParallelism(1)
.print();
env.execute("Flink CDC MySQL to Console");
注意事项:
- 表必须有主键,否则无法生成 update/delete 事件。
- 需提前在 MySQL 创建用户并授权:
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flink'@'%';
- 输出为 JSON 格式的变更事件,包含
op(操作类型)、ts_ms(时间戳)、before、after等字段。
第3章 Flink CDC 核心 API 详解
3.1 MySqlSource 构建器模式详解
| 方法名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
builder() | MySqlSource.<T>builder() | 获取构建器实例,泛型 T 为输出数据类型 | MySqlSource.<String>builder() | 必须调用 build() 前使用 |
hostname() | .hostname(String hostname) | 设置 MySQL 服务器地址 | .hostname("192.168.1.100") | 不支持端口拼接(如 host:port) |
port() | .port(int port) | 设置 MySQL 端口 | .port(3306) | 默认值为 3306 |
username() | .username(String username) | 设置连接用户名 | .username("flink_user") | 需具备 REPLICATION 权限 |
password() | .password(String password) | 设置连接密码 | .password("SecurePass123!") | 建议使用配置文件或密钥管理 |
databaseName() | .databaseName(String databaseName) | 指定监听的数据库名(支持正则) | .databaseName("prod_db") / .databaseName("test_.*") | 正则需用引号包裹 |
tableName() | .tableName(String tableName) | 指定监听的表名(支持正则) | .tableName("users") / .tableName("order_.*") | 支持多表匹配 |
serverId() | .serverId(String serverId) | 设置 MySQL server-id(用于 binlog 读取) | .serverId("5075") | 推荐设置为唯一整数,格式为 "startId-endId"(如 "5050-5080") |
serverTimeZone() | .serverTimeZone(String timeZone) | 设置服务器时区 | .serverTimeZone("Asia/Shanghai") | 影响时间字段解析,避免时区偏移 |
startupMode() | .startupMode(StartupMode mode) | 设置启动模式 | .startupMode(StartupMode.INITIAL) | 见 3.4 节详解 |
deserializer() | .deserializer(DeserializationSchema<T> deserializer) | 设置反序列化器 | .deserializer(ChangelogJsonDeserializationSchema.INSTANCE) | 决定输出数据格式 |
debeziumConfig() | .debeziumConfig(String key, String value) | 添加 Debezium 底层配置 | .debeziumConfig("snapshot.mode", "schema_only") | 可覆盖默认行为 |
includeSchemaChanges() | .includeSchemaChanges(boolean include) | 是否包含 DDL 事件 | .includeSchemaChanges(true) | 实验性功能,需目标系统支持 |
splitSize() | .splitSize(Integer splitSize) | 设置快照分片大小(行数) | .splitSize(8000) | 控制并行读取粒度,影响性能 |
fetchSize() | .fetchSize(Integer fetchSize) | 设置每次读取的记录数 | .fetchSize(1024) | 提升网络传输效率 |
connectTimeout() | .connectTimeout(Duration timeout) | 连接超时时间 | .connectTimeout(Duration.ofSeconds(30)) | 防止长时间阻塞 |
connectMaxRetries() | .connectMaxRetries(Integer maxRetries) | 连接重试次数 | .connectMaxRetries(3) | 应对网络抖动 |
build() | .build() | 构建 MySqlSource 实例 | MySqlSource<String> source = builder.build(); | 必须调用,否则无法创建 Source |
注: 所有方法均为链式调用,最终调用
build()返回SourceFunction<T>。
3.2 数据消费方式:SourceFunction 与 DataStream 集成
| 方法/接口 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
addSource() | env.addSource(source) | 将 CDC Source 添加到流环境中 | DataStreamSource ds = env.addSource(mySqlSource); | 是接入 CDC 的入口方法 |
setParallelism() | .setParallelism(n) | 设置算子并行度 | ds.setParallelism(1) | MySqlSource 通常并行度为 1(单分片读取 binlog) |
print() | .print() | 输出到控制台(调试用) | ds.print(); | 输出带并行子任务编号 |
map() / filter() / flatMap() | .map(...), .filter(...), .flatMap(...) | 转换或过滤变更数据 | ds.map(json -> JSON.parseObject(json)) | 常用于解析 JSON 或投影字段 |
keyBy() | .keyBy(...) | 按字段分组(用于窗口计算) | ds.keyBy(json -> JSON.parseObject(json).getString("after.id")) | 需确保字段存在 |
addSink() | .addSink(sinkFunction) | 自定义 Sink 输出 | ds.addSink(new JdbcSink()); | 可实现数据库写入等操作 |
execute() | env.execute("Job Name") | 触发作业执行 | env.execute("MySQL CDC Job"); | 必须调用,否则不运行 |
完整集成示例:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
MySqlSource source = MySqlSource.builder()
.hostname("localhost")
.databaseName("test")
.tableName("user")
.username("flink")
.password("cdc")
.deserializer(ChangelogJsonDeserializationSchema.INSTANCE)
.build();
DataStream stream = env.addSource(source).setParallelism(1);
stream.map(record -> parseUserFromJson(record))
.addSink(new CustomUserSink());
env.execute("User Sync Job");
注意事项:
- MySqlSource 默认并行度为 1,不支持并行读取 binlog。
- 若需并行处理,可在
map等后续算子中提高并行度。 SourceFunction是 Flink 原生接口,兼容所有版本。
3.3 Debezium 配置参数集成与调优
| 参数名 | 语法示例 | 用途 | 推荐值 | 注意事项 |
|---|---|---|---|---|
snapshot.mode | .debeziumConfig("snapshot.mode", "initial") | 控制快照行为 | initial, schema_only, schema_only_recovery | initial=全量+增量,schema_only=仅结构 |
snapshot.locking.mode | .debeziumConfig("snapshot.locking.mode", "none") | 快照锁策略 | none, minimal, extended | none 表示无锁读取(推荐) |
binlog.buffer.size | .debeziumConfig("binlog.buffer.size", "16384") | binlog 缓冲区大小 | 16384 ~ 1048576 | 单位字节,提升吞吐 |
connect.timeout.ms | .debeziumConfig("connect.timeout.ms", "30000") | 连接超时 | 30000(30秒) | 防止长时间阻塞 |
socket.timeout.ms | .debeziumConfig("socket.timeout.ms", "60000") | Socket 读取超时 | 60000(60秒) | 应对网络延迟 |
poll.interval.ms | .debeziumConfig("poll.interval.ms", "100") | 轮询间隔 | 100 ~ 500 | 值越小延迟越低,CPU 消耗越高 |
decimal.handling.mode | .debeziumConfig("decimal.handling.mode", "string") | DECIMAL 字段处理方式 | precise(保留精度), string | string 更安全 |
skip.messages | .debeziumConfig("skip.messages", "true") | 跳过无法解析的消息 | true / false | 避免作业失败 |
heartbeat.interval.ms | .debeziumConfig("heartbeat.interval.ms", "30000") | 心跳发送间隔 | 30000(30秒) | 用于监控连接状态 |
调优建议:
- 生产环境建议设置
snapshot.locking.mode=none避免锁表。 - 大表同步可调大
binlog.buffer.size和poll.interval.ms降低 CPU。 - 使用
decimal.handling.mode=string防止精度丢失。
3.4 启动模式(StartupMode)详解:initial、latest-offset、timestamp、specific-offset
| 启动模式 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
INITIAL | StartupMode.INITIAL | 首次启动:先全量快照,再读增量 binlog | .startupMode(StartupMode.INITIAL) | 适用于新表同步 |
LATEST_OFFSET | StartupMode.LATEST_OFFSET | 仅从当前最新位点开始读取 binlog | .startupMode(StartupMode.LATEST_OFFSET) | 不读历史数据,常用于灾备恢复 |
TIMESTAMP | StartupMode.TIMESTAMP | 从指定时间戳开始读取 | .startupMode(StartupMode.TIMESTAMP) / .startupTimestamp(1672531200L) | 时间戳单位为秒 |
SPECIFIC_OFFSET | StartupMode.SPECIFIC_OFFSET | 从指定 binlog 文件和位置开始 | .startupMode(StartupMode.SPECIFIC_OFFSET) / .startupSpecificOffset("mysql-bin.000003", 154L) | 精准恢复场景 |
NEVER | StartupMode.NEVER | 仅读取未来变更(无快照) | .startupMode(StartupMode.NEVER) | 表必须已存在且不关心历史数据 |
使用场景说明:
INITIAL:最常用,确保数据完整。LATEST_OFFSET:跳过历史数据,快速接入。TIMESTAMP:按时间恢复,如”从昨天开始同步”。SPECIFIC_OFFSET:故障恢复时指定精确位点。
3.5 输出模式(OutputMode)详解:debezium、canal、changelog-json
| 输出模式 | 对应反序列化器 | 用途 | 输出示例片段 | 注意事项 |
|---|---|---|---|---|
| Debezium 格式 | DebeziumSourceFunction | 原生 Debezium JSON 格式 | {"op":"c","ts_ms":...,"before":null,"after":{"id":1,"name":"Alice"}} | 字段丰富,适合 Kafka 中转 |
| Canal 格式 | CanalJsonDeserializationSchema | 阿里开源的 Canal 协议格式 | {"type":"INSERT","ts":...,"data":[{"id":"1","name":"Bob"}]} | 兼容 Canal 生态 |
| Changelog JSON | ChangelogJsonDeserializationSchema | Flink 自定义格式,兼容 changelog stream | ["+I",{"id":1,"name":"Charlie"}] | "+I"=insert, "-U"=旧值 update, "+U"=新值 update, "-D"=delete |
| 自定义 RowData | RowDataDeserializationSchema | 输出为 Flink 内部 RowData 类型 | RowData with binary fields | 高性能,适合与 Flink SQL 集成 |
代码示例:
// 使用 Changelog JSON
.deserializer(ChangelogJsonDeserializationSchema.INSTANCE)
// 使用 Canal JSON
.deserializer(CanalJsonDeserializationSchema.builder().build())
// 使用 Debezium JSON
.deserializer(new JsonDebeziumDeserializationSchema())
注意事项:
ChangelogJsonDeserializationSchema输出为 String,需进一步解析。RowData模式需配合 Flink Table API 使用,性能最优。- 不同格式字段结构不同,下游需适配解析逻辑。
第4章 变更事件处理与数据解析
4.1 CDC 数据结构解析(Debezium 格式)
| 字段名 | 说明 | 示例值 | 注意事项 |
|---|---|---|---|
op | 操作类型 | "c"=create(insert), "u"=update, "d"=delete, "r"=read(快照) | "r" 仅出现在快照阶段 |
ts_ms | 事件发生时间(毫秒) | 1672531200000 | 来自数据库时间,受时区影响 |
before | 变更前的数据(update/delete) | {"id":1,"name":"Alice"} | insert 时为 null |
after | 变更后的数据(insert/update) | {"id":2,"name":"Bob"} | delete 时为 null |
source | 源元数据 | {"version":"1.9.7.Final", "db":"test_db", "table":"users"} | 包含库、表、事务ID等 |
transaction | 事务信息(如有) | {"id":"123", "total_order":5} | 需开启事务支持 |
databaseName | 数据库名(顶层字段) | "test_db" | 便于路由不同库的表 |
完整事件结构示例(JSON):
{
"op": "u",
"ts_ms": 1672531200123,
"before": {"id": 1, "name": "Alice", "age": 25},
"after": {"id": 1, "name": "Alice", "age": 26},
"source": {
"db": "test_db",
"table": "users",
"server_id": 5075
}
}
注意事项:
before和after是嵌套对象,需递归解析。- 时间字段可能为字符串或时间戳,取决于配置。
- 大字段(如 BLOB)可能被截断或 Base64 编码。
4.2 如何解析 insert/update/delete 事件
| 事件类型 | 判断条件 | 处理逻辑 | 代码示例(Java) | 注意事项 |
|---|---|---|---|---|
| Insert | op == "c" 或 "+I" | 取 after 字段作为新数据 | if (op.equals("c")) { User user = parseAfter(json); } | 快照中的 insert 也标记为 "c" |
| Update | op == "u" | before 为旧值,after 为新值 | if (op.equals("u")) { User old = parseBefore(json); User updated = parseAfter(json); } | 区分旧值和新值 |
| Delete | op == "d" | 取 before 字段作为被删数据 | if (op.equals("d")) { User deleted = parseBefore(json); } | after 为 null |
| Changelog JSON | 记录首字段 | ["+I", data], ["-U", old], ["+U", new], ["-D", data] | String op = array.getString(0); if ("+I".equals(op)) { ... } | 更简洁,推荐用于 Flink 内部处理 |
通用解析逻辑:
JsonObject obj = JsonParser.parse(json).getAsJsonObject();
String op = obj.get("op").getAsString();
switch (op) {
case "c": // insert
handleInsert(obj.getAsJsonObject("after"));
break;
case "u": // update
handleUpdate(obj.getAsJsonObject("before"), obj.getAsJsonObject("after"));
break;
case "d": // delete
handleDelete(obj.getAsJsonObject("before"));
break;
}
4.3 使用 RowData 与 JSONObject 处理原始变更
| 方法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|
| JSONObject(JSON 字符串) | 解析 Debezium/Canal 输出的 JSON | JsonObject json = JsonParser.parse(text).getAsJsonObject(); / String op = json.get("op").getAsString(); | 需引入 fastjson 或 gson |
| RowData(Flink 内部类型) | 高性能二进制格式,与 Table API 无缝集成 | RowData rowData = deserializer.deserialize(sourceRecord); / String op = rowData.getRowKind().shortString(); | 需使用 RowDataDeserializationSchema |
| RowKind | 表示变更类型 | if (rowData.getRowKind() == RowKind.INSERT) { ... } | INSERT, UPDATE_BEFORE, UPDATE_AFTER, DELETE |
| BinaryConverter | 将 RowData 字段转为 Java 类型 | String name = BinaryStringData.toString(rowData.getString(1)); | 需注意字段类型和索引 |
RowData 示例:
MySqlSource source = MySqlSource.builder()
.deserializer(RowDataDeserializationSchema.builder()
.build())
.build();
DataStream stream = env.addSource(source);
stream.map(row -> {
RowKind kind = row.getRowKind();
if (kind == RowKind.INSERT) {
// 处理插入
}
});
4.4 自定义反序列化器(DeserializationSchema)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
deserialize() | T deserialize(byte[] message) | 核心方法:将字节数组转为对象 | public User deserialize(byte[] msg) { return JSON.parseObject(new String(msg), User.class); } | 抛出 IOException 表示解析失败 |
isEndOfStream() | boolean isEndOfStream(T nextElement) | 是否终止流 | return false; | 通常返回 false(无限流) |
getProducedType() | TypeInformation getProducedType() | 声明输出类型 | return TypeInformation.of(User.class); | 供 Flink 类型系统使用 |
open() | void open(InitializationContext context) | 初始化资源(如连接池) | public void open() { objectMapper = new ObjectMapper(); } | 可选,Task 初始化时调用 |
完整自定义反序列化器示例:
public class UserDeserializationSchema implements DeserializationSchema {
private ObjectMapper mapper;
@Override
public void open(InitializationContext context) {
mapper = new ObjectMapper();
}
@Override
public User deserialize(byte[] message) throws IOException {
JsonNode node = mapper.readTree(message);
JsonNode after = node.get("after");
return mapper.treeToValue(after, User.class);
}
@Override
public boolean isEndOfStream(User nextElement) {
return false;
}
@Override
public TypeInformation getProducedType() {
return TypeInformation.of(User.class);
}
}
注意事项:
- 必须实现
DeserializationSchema<T>接口。 deserialize方法必须高效,避免阻塞。- 使用
open()初始化对象(如 ObjectMapper),避免重复创建。
第5章 水印与事件时间处理
5.1 CDC 场景下的事件时间语义支持
| 概念名称 | 说明 | 用途 | 注意事项 |
|---|---|---|---|
| 事件时间(Event Time) | 数据在源头数据库中发生变更的时间,通常取自 ts_ms 字段。 | 用于窗口计算、乱序处理,保证结果一致性。 | 必须从 CDC 消息中提取时间字段,不能使用系统处理时间。 |
| 处理时间(Processing Time) | Flink 算子接收到数据的时间。 | 简单实时处理,如监控告警。 | 不保证结果一致性,不适用于精确统计。 |
| 摄取时间(Ingestion Time) | 数据进入 Flink Source 算子的时间。 | 折中方案,延迟较低且有一定有序性。 | 仍可能受 Source 内部排队影响。 |
ts_ms 字段 | Debezium 输出中的 ts_ms,表示事件在数据库中的发生时间(毫秒)。 | 作为事件时间的时间戳来源。 | 来自数据库服务器时间,需确保时钟同步。 |
source.ts_ms | 更精确的字段,表示 binlog 写入时间(MySQL 5.7+)。 | 比 ts_ms 更接近真实变更时间。 | 推荐优先使用 source.ts_ms。 |
说明: 在 CDC 场景中,应优先使用事件时间,以确保窗口聚合等操作的准确性,尤其是在数据延迟或重放时。
5.2 基于变更记录生成水印(WatermarkStrategy)
| 方法/类 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
forBoundedOutOfOrderness() | WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)) | 适用于有界乱序场景,允许最大延迟 N 秒 | WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(10)) | 最常用策略,需合理设置延迟时间 |
forMonotonousTimestamps() | .forMonotonousTimestamps() | 时间戳单调递增(极少乱序) | .forMonotonousTimestamps() | 不适用于真实 CDC 场景(可能有延迟) |
noWatermarks() | .noWatermarks() | 不生成水印,仅用于测试 | WatermarkStrategy.noWatermarks() | 窗口无法触发,生产禁用 |
withTimestampAssigner() | .withTimestampAssigner((event, timestamp) -> {...}) | 提取事件时间戳 | 见下例 | 必须返回 long 类型时间戳(毫秒) |
| Watermark | 水印机制 | 标记事件时间的进展,触发窗口计算 | Flink 自动管理 | 水印 < 所有未处理事件的时间戳 - 延迟 |
完整代码示例:
DataStream stream = env.addSource(mySqlSource);
WatermarkStrategy strategy = WatermarkStrategy
.forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((event, timestamp) -> {
try {
JsonObject obj = JsonParser.parse(event).getAsJsonObject();
// 优先使用 source.ts_ms
JsonElement source = obj.get("source");
if (source != null && source.isJsonObject()) {
return source.getAsJsonObject().get("ts_ms").getAsLong();
}
return obj.get("ts_ms").getAsLong();
} catch (Exception e) {
return timestamp; // 解析失败使用系统时间
}
});
stream.assignTimestampsAndWatermarks(strategy)
.keyBy(json -> extractUserId(json))
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new UserCountAgg())
.print();
注意事项:
- 水印生成必须在
keyBy和window之前调用。 - 延迟时间(bounded delay)应根据业务容忍度和网络延迟设置。
- 若时间字段为空或异常,应提供默认值或日志告警。
5.3 处理延迟数据与乱序事件
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
allowedLateness() | .allowedLateness(Time.minutes(5)) | 允许延迟数据在窗口关闭后继续到达 | window(TumblingEventTimeWindows.of(Time.minutes(10))).allowedLateness(Time.minutes(2)) | 延迟数据会触发窗口再次计算 |
sideOutputLateData() | OutputTag lateTag = new OutputTag<>("late-data"){}; ... .sideOutputLateData(lateTag) | 将迟到数据输出到侧输出流 | 见下方示例 | 可用于告警或异步处理 |
| late element | 延迟元素 | 时间戳 < 当前水印 - 延迟阈值的数据 | 自动被丢弃或路由到侧输出流 | 应监控其数量,判断系统健康度 |
| Trigger | 触发器 | 控制窗口何时计算 | 默认为 EventTimeTrigger | 可自定义触发逻辑(如连续5条数据触发) |
处理策略对比:
| 策略 | 适用场景 | 实现方式 | 优缺点 |
|---|---|---|---|
| 丢弃 | 对延迟不敏感 | 默认行为 | 简单,但可能丢失数据 |
| 允许迟到 | 可容忍短时延迟 | allowedLateness() | 提高准确性,增加状态开销 |
| 侧输出流 | 需单独处理延迟数据 | sideOutputLateData() | 灵活,可重试或告警,需额外处理逻辑 |
侧输出流示例:
OutputTag lateTag = new OutputTag<>("late-data"){};
SingleOutputStreamOperator mainStream = windowedStream
.sideOutputLateData(lateTag)
.aggregate(new AggFunc());
DataStream lateStream = mainStream.getSideOutput(lateTag);
建议:
- 对关键指标(如交易额),设置
allowedLateness(1-5分钟)。 - 对非关键数据,使用侧输出流记录日志。
- 监控
numLateRecordsDropped指标,评估延迟情况。
第6章 容错与一致性保障
6.1 Checkpointing 机制在 CDC 中的作用
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
enableCheckpointing() | env.enableCheckpointing(5000) | 启用检查点,每5秒一次 | env.enableCheckpointing(5000); | 是实现 Exactly-Once 的基础 |
setCheckpointingMode() | .setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE) | 设置检查点模式 | .setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE) | 推荐使用 EXACTLY_ONCE |
setCheckpointTimeout() | .setCheckpointTimeout(60000) | 检查点超时时间 | .setCheckpointTimeout(60000) | 超时后会被丢弃 |
setMinPauseBetweenCheckpoints() | .setMinPauseBetweenCheckpoints(500) | 两次检查点最小间隔 | .setMinPauseBetweenCheckpoints(500) | 防止频繁触发 |
setMaxConcurrentCheckpoints() | .setMaxConcurrentCheckpoints(1) | 最大并发检查点数 | .setMaxConcurrentCheckpoints(1) | 通常为1,避免资源竞争 |
enableExternalizedCheckpoints() | .enableExternalizedCheckpoints(...) | 外部化保存检查点 | .enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION) | 作业取消后保留检查点 |
Checkpoint 在 CDC 中的关键作用:
- 记录当前读取的 binlog 位点(filename + position)。
- 保证 Source、Transformation、Sink 的状态一致性。
- 故障恢复时,从最近检查点恢复,避免数据丢失或重复。
代码示例:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(5000); // 每5秒一次
CheckpointConfig config = env.getCheckpointConfig();
config.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
config.setCheckpointTimeout(60000);
config.setMinPauseBetweenCheckpoints(500);
config.setMaxConcurrentCheckpoints(1);
config.enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
6.2 Exactly-Once 语义实现原理
| 组件 | 作用 | 实现机制 | 注意事项 |
|---|---|---|---|
| Flink Checkpoint | 分布式快照 | 基于 Chandy-Lamport 算法,全局一致性快照 | 需 Barrier 对齐(在 EXACTLY_ONCE 模式下) |
| Source(CDC) | 可回放的流 | 记录 binlog 位点,支持从指定位置重读 | MySQL binlog 必须保留足够长时间 |
| Sink | 幂等写入或事务提交 | - 幂等:如 Redis SET key value - 事务:如 Kafka 事务、Doris Stream Load 事务 | 推荐使用事务型 Sink |
| Two-Phase Commit (2PC) | 预提交与提交 | Flink 提供 TwoPhaseCommitSinkFunction | 是 Exactly-Once 的关键保障 |
| Barrier | 检查点分界符 | 在数据流中插入特殊标记 | 所有算子需对其对齐 |
Exactly-Once 流程:
- Flink 触发 Checkpoint。
- Source 将当前 binlog 位点作为状态保存。
- Barrier 流经所有算子,触发状态快照。
- Sink 执行预提交(pre-commit),保存事务状态。
- JobManager 确认所有任务快照成功。
- Sink 提交事务(commit),释放资源。
注意事项:
- 若 Sink 不支持事务或幂等,无法保证端到端 Exactly-Once。
- 网络分区或任务失败时,Flink 会回滚到上一个检查点重新消费。
6.3 故障恢复与 Binlog 位点自动恢复
| 恢复机制 | 说明 | 触发条件 | 注意事项 |
|---|---|---|---|
| 自动从 Checkpoint 恢复 | Flink 重启时自动读取最近的检查点 | 作业失败、手动重启 | 需启用检查点和外部化存储 |
| Binlog 位点恢复 | Source 根据检查点中保存的 filename 和 position 重新连接 MySQL | 启动模式为 INITIAL 或 TIMESTAMP 时 | 要求 binlog 文件未被清理 |
| 断点续传 | 类似下载断点续传,继续上次中断的位置 | 网络抖动、临时故障 | 依赖 Debezium 的 offset 存储机制 |
| externalized checkpoints | 外部化检查点(如 HDFS、S3) | 作业取消后仍保留 | 可用于版本升级或迁移 |
| savepoint | 手动生成的检查点 | flink savepoint 命令 | 用于版本升级、A/B 测试 |
恢复流程:
- 作业重启。
- Flink 从持久化存储加载最近的 Checkpoint/Savepoint。
- MySqlSource 读取其中的 binlog 位点(offset)。
- 连接 MySQL,从该位点继续读取 binlog。
- 继续处理后续变更。
注意事项:
- 必须确保 MySQL 的
expire_logs_days或binlog_expire_logs_seconds足够长,覆盖最大恢复时间。 - 若 binlog 已被清理,将导致恢复失败,需重新全量同步。
6.4 并行度与分片策略对一致性的影响
| 策略 | 说明 | 一致性影响 | 注意事项 |
|---|---|---|---|
| MySqlSource 并行度=1 | 单任务读取 binlog | 强一致性,事件顺序完全保序 | 吞吐受限,适用于中小数据量 |
| 并行度>1 | Flink CDC 当前不支持 | 若强行设置,可能导致位点混乱 | 禁止手动设置并行度 > 1 |
| 分片快照(Split Snapshot) | 大表快照时自动分片读取 | 快照阶段无全局一致性(非瞬时快照) | 通过 split.size 控制分片大小 |
| 全局一致性快照 | 所有表在同一时刻的快照 | 需使用 RELOAD 权限和 FLUSH TABLES WITH READ LOCK | 影响数据库性能,不推荐 |
| 无锁快照(Lock-free) | 基于 RR 隔离级别和主键范围扫描 | 快照期间允许写入,但保证最终一致性 | 推荐使用,对业务影响小 |
| 多表同步 | 单 Job 监听多个表 | 各表快照时间不同 | 无法保证跨表事务一致性 |
最佳实践:
- MySqlSource 的并行度必须为 1,以保证 binlog 读取顺序。
- 大表使用
split.size启用分片快照,提升性能。 - 若需高吞吐,可通过部署多个 Job 分散不同表的同步压力。
- 跨表一致性需在业务层或 Sink 层处理(如使用事务型数据库)。
第7章 多表同步与元数据管理
7.1 单 Job 同步多个表(正则匹配表名)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
tableName() | .tableName(String regex) | 使用正则表达式匹配多个表 | .tableName("user_.*") / .tableName("(user|order|product)") | 正则需用引号包裹 |
databaseName() | .databaseName(String regex) | 匹配多个数据库中的表 | .databaseName("prod_.*") | 可与 tableName() 组合使用 |
| 多模式组合 | — | 实现库/表级过滤 | .databaseName("finance|hr").tableName("(user|dept)_.*") | 使用 | 分隔多个模式 |
| 白名单机制 | 内置支持 | 基于正则实现白名单过滤 | 见上例 | 不支持黑名单(blacklist),需在下游过滤 |
完整示例:
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseName("test_.*") // 所有 test_ 开头的库
.tableName("(user|order|log)_.*") // 用户、订单、日志类表
.username("flink")
.password("cdc")
.deserializer(ChangelogJsonDeserializationSchema.INSTANCE)
.build();
优势: 减少 Job 数量,统一运维。
注意: 所有匹配表必须结构兼容反序列化器输出类型;若结构差异大,建议拆分 Job 或使用通用格式(如 JSON)。
7.2 表结构元数据获取(schema changes)
| 字段/功能 | 说明 | 示例值 | 注意事项 |
|---|---|---|---|
source.struct.version | Debezium 结构版本 | "v2" | 标识 schema 描述格式 |
source.schema | 表结构定义(DDL 信息) | { "type": "struct", "fields": [ ... ] } | 包含字段名、类型、是否主键等 |
before / after 结构变化 | DDL 后字段增删改 | 新增字段出现在 after 中 | before 可能缺失新字段 |
op = "c" + ddl 字段 | DDL 操作事件 | "op":"c", "ddl":"ALTER TABLE users ADD COLUMN email VARCHAR(255)" | 仅当 .includeSchemaChanges(true) 时输出 |
includeSchemaChanges() | 是否包含 DDL 事件 | .includeSchemaChanges(true) | 实验性功能,部分 Sink 需适配处理 DDL |
| 数据血缘(Data Lineage) | 基于 schema 变更构建 | 记录字段生命周期 | 可用于元数据中心建设 |
启用 DDL 监听示例:
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.databaseName("test")
.tableName("users")
.includeSchemaChanges(true) // 启用 DDL 输出
.deserializer(ChangelogJsonDeserializationSchema.INSTANCE)
.build();
典型 DDL 事件 JSON:
{
"op": "c",
"ts_ms": 1712345678901,
"ddl": "ALTER TABLE users ADD COLUMN email VARCHAR(255) AFTER name",
"source": { "db": "test", "table": "users" }
}
限制:
- Flink CDC 当前不支持自动响应 DDL(如动态修改 RowData Schema)。
- 需在下游手动处理或忽略 DDL 事件。
7.3 动态添加监听表(Future Feature & Workaround)
| 方案 | 说明 | 实现方式 | 优缺点 | 注意事项 |
|---|---|---|---|---|
| 原生不支持 | Flink CDC 目前无法运行时动态增表 | 构建时固定表列表 | 缺陷:需重启 Job | 官方 roadmap 中规划中 |
| 多 Job 分治 | 每个表或每组表独立 Job | 部署 N 个 Job,按命名规则管理 | 灵活,易扩展 / 运维成本高 | 推荐用于生产环境 |
| 控制台配置 + 重启 | 通过外部配置中心(如 Nacos)管理表名正则,更新后滚动重启 Job | 配置变更 → CI/CD 自动部署 | 半动态 / 有短暂中断 | 适用于低频变更场景 |
| Kafka Connect + Debezium | 使用 Kafka Connect 框架替代 Flink CDC | Connect 支持动态发现新表 | 真正动态 / 脱离 Flink 生态 | 适合已使用 Kafka 的架构 |
| Flink Application Mode + Savepoint | 修改 Job 逻辑并从 Savepoint 恢复 | flink run -s <savepoint> ... | 保证状态恢复 / 需重新编译打包 | 适合版本升级 |
建议: 优先采用”多 Job + 配置化”方案,平衡灵活性与稳定性。
7.4 数据路由:将不同表写入不同 Sink
| 方法 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
sideOutput + ProcessFunction | 使用侧输出流分流 | 见下方完整示例 | 最灵活,推荐使用 |
split() (Deprecated) | 已废弃,不推荐 | — | 避免使用 |
| 多 DataStream 分支 | map/filter 后分别 addSink | usersStream.addSink(userSink); / ordersStream.addSink(orderSink); | 适用于简单路由 |
| 自定义 Sink 路由逻辑 | 在 Sink 内部判断表名并路由 | if (table.equals("user")) writeMySQL(); else writeES(); | 减少 Job 图复杂度 |
完整路由示例(使用 ProcessFunction + Side Output):
// 定义侧输出标签
OutputTag<String> userTag = new OutputTag<>("user-data"){};
OutputTag<String> orderTag = new OutputTag<>("order-data"){};
OutputTag<String> logTag = new OutputTag<>("log-data"){};
SingleOutputStreamOperator<String> routed = stream
.process(new ProcessFunction<String, String>() {
@Override
public void processElement(String json, Context ctx, Collector<String> out) {
JsonObject obj = JsonParser.parse(json).getAsJsonObject();
String table = obj.getAsJsonObject("source").get("table").getAsString();
if (table.startsWith("user")) {
ctx.output(userTag, json);
} else if (table.startsWith("order")) {
ctx.output(orderTag, json);
} else if (table.startsWith("log")) {
ctx.output(logTag, json);
} else {
out.collect(json); // 主流
}
}
});
// 获取各侧输出流并写入不同 Sink
DataStream<String> userStream = routed.getSideOutput(userTag);
DataStream<String> orderStream = routed.getSideOutput(orderTag);
DataStream<String> logStream = routed.getSideOutput(logTag);
userStream.addSink(new JdbcSink("jdbc:mysql://.../dw_users"));
orderStream.addSink(new KafkaSink("order_topic"));
logStream.addSink(new ElasticsearchSink());
优点: 单 Job 实现多目的地写入,资源利用率高。
注意: 确保各 Sink 的失败策略一致,避免部分成功导致状态不一致。
第8章 实时数仓集成实践
8.1 CDC → Kafka:作为消息中间件中转
| 配置项 | 说明 | 示例 | 注意事项 |
|---|---|---|---|
| Sink 类型 | Flink Kafka Producer | FlinkKafkaProducer<String> | 需引入 flink-connector-kafka |
| 序列化格式 | 建议使用 JSON 或 AVRO | Changelog JSON、Debezium JSON | 便于下游消费 |
| Key 设计 | 使用主键哈希提升并行度 | .setKeyedSerializationSchema(...) | 避免热点 |
| 事务写入 | 保障 Exactly-Once | .setWriteTransactionsTimeout(...) | 需 Kafka 0.11+ |
| Topic 路由 | 按表名映射到不同 Topic | db.table.users, db.table.orders | 命名规范便于管理 |
代码示例:
KafkaSink<String> kafkaSink = KafkaSink.<String>builder()
.setBootstrapServers("kafka:9092")
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("cdc-changelog") // 可动态设置 topic
.setValueSerializationSchema(SimpleStringSchema.INSTANCE)
.build())
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setProperty("transaction.timeout.ms", "60000")
.build();
stream.sinkTo(kafkaSink);
用途: 解耦数据生产与消费,支持多订阅者(如 OLAP、数仓、审计)。
安全: 建议启用 SSL/SASL 认证。
8.2 CDC → JDBC:写入其他数据库(MySQL、PostgreSQL)
| 参数 | 说明 | 示例 | 注意事项 |
|---|---|---|---|
| Sink 类型 | JdbcSink | JdbcSink.sink(...) | 支持批处理写入 |
| 批量大小 | 控制写入频率 | .withBatchSize(1000) | 平衡延迟与吞吐 |
| 批量间隔 | 超时强制提交 | .withBatchIntervalMs(2000) | 防止数据积压 |
| 更新策略 | insert or update | ON DUPLICATE KEY UPDATE(MySQL)/ ON CONFLICT DO UPDATE(PG) | 实现 upsert |
| SQL 模板 | 动态生成 SQL | "INSERT INTO users VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE name=VALUES(name)" | 支持参数占位符 |
Upsert 示例(MySQL):
JdbcConnectionOptions options = new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:mysql://target:3306/dw")
.withUsername("dw")
.withPassword("pass")
.build();
JdbcSink<String> sink = JdbcSink.<String>sink(
"INSERT INTO users(id, name, age) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE name=VALUES(name), age=VALUES(age)",
(ps, json) -> {
JsonObject after = JsonParser.parse(json).getAsJsonObject().get("after").getAsJsonObject();
ps.setInt(1, after.get("id").getAsInt());
ps.setString(2, after.get("name").getAsString());
ps.setInt(3, after.get("age").getAsInt());
},
JdbcExecutionOptions.builder()
.withBatchSize(1000)
.withBatchIntervalMs(2000)
.build(),
options
);
注意:
- 目标表需有主键才能实现 upsert。
- 高频更新场景建议使用 Doris/ClickHouse 替代传统 RDBMS。
8.3 CDC → Hudi / Iceberg:构建湖仓一体架构
| 特性 | Hudi | Iceberg | 说明 |
|---|---|---|---|
| 写入模式 | COPY_ON_WRITE / MERGE_ON_READ | APPEND / UPSERT | 支持 CDC 增量写入 |
| Flink 集成 | HoodieFlinkStreamer | FlinkIcebergSink | 均支持流式写入 |
| 更新支持 | 支持 update/delete | 支持 changelog | 适合 CDC 场景 |
| 时间旅行 | 支持 | 支持 | 基于 snapshot 查询历史状态 |
| 元数据管理 | 内嵌文件列表 | 元数据文件(Manifest) | Iceberg 更轻量 |
| 小文件合并 | 自动压缩 | 需手动 compact | 影响查询性能 |
Hudi 写入示例(配置方式):
# flink-conf.yaml
execution.checkpointing.interval: 5min
-- SQL 定义 Hudi 表
CREATE TABLE hudi_users (
id INT PRIMARY KEY,
name STRING,
age INT,
ts TIMESTAMP(3),
`__changelog` STRING
) PARTITIONED BY (dt)
WITH (
'connector' = 'hudi',
'path' = 's3a://lake/users',
'table.type' = 'MERGE_ON_READ',
'write.precombine.field' = 'ts',
'write.operation' = 'upsert'
);
优势: 支持 ACID、增量查询、与 Hive/Spark/Presto 兼容。
适用场景: 实时数仓、数据湖、机器学习特征存储。
8.4 CDC → Doris / ClickHouse:实时 OLAP 分析
| 对比项 | Doris | ClickHouse |
|---|---|---|
| 写入方式 | Stream Load(HTTP) | HTTP Interface / Native TCP |
| Flink Connector | DorisSink | ClickHouseSink(社区) |
| UPSERT 支持 | 唯一模型表 | CollapsingMergeTree / ReplacingMergeTree |
| 实时性 | 秒级 | 秒级 |
| 语法兼容 | MySQL | SQL-like |
| 适用场景 | 实时报表、BI 分析 | 高并发日志分析、用户行为 |
Doris Sink 示例:
DorisSink.Builder<String> builder = DorisSink.builder();
builder.setDorisConfig(Maps.of(
"fenodes", "doris-fe:8030",
"username", "admin",
"password", "",
"table.identifier", "analytics.users"
));
builder.setSerializer(new SimpleStringSchema()); // 或自定义序列化
DorisSink<String> dorisSink = builder.build();
stream.sinkTo(dorisSink);
ClickHouse Sink 示例(使用 ReplacingMergeTree):
-- 表结构需包含 version 字段
CREATE TABLE users (
id UInt32,
name String,
age UInt8,
version UInt64
) ENGINE = ReplacingMergeTree(version)
ORDER BY id;
优势: 高吞吐写入、低延迟查询,适合构建实时大屏、用户画像。
注意: 需合理设计主键和排序键,避免性能瓶颈。
第9章 性能调优与生产最佳实践
9.1 并行读取配置(split-size、fetch-size)
注: Flink CDC 的 binlog 读取阶段为单并行度,但**快照(snapshot)**阶段支持分片并行读取。本节聚焦快照性能优化。
| 配置项 | 默认值 | 说明 | 调优建议 | 注意事项 |
|---|---|---|---|---|
scan.snapshot.fetch.size | 1024 | 每次从数据库 fetch 的行数 | 增大至 512~4096 提升吞吐 | 受 JDBC fetchSize 限制,过大可能占用内存 |
scan.split.size | 8096 | 每个分片(split)的行数 | 大表设为 10000~50000,提升并行度 | 值越小,并行任务越多,但小文件多 |
scan.incremental.snapshot.chunk.size | 1024 | 增量快照分块大小(行) | 控制 binlog 回放粒度 | 与 split.size 类似,用于增量阶段 |
parallelism | 1(Source) | 快照阶段并行度 | 设置 env.setParallelism(N) | 实际并行度 = min(N, 分片总数) |
scan.split.mode | "primary-key" | 分片依据:主键范围扫描 | 可选 "range" 或 "hybrid" | 主键需为数字或时间类型 |
调优示例:
MySqlSource.<String>builder()
.databaseName("prod")
.tableName("large_table")
.scanStartupMode(StartupMode.INITIAL) // 全量 + 增量
.splitSize(20000) // 每个 split 2 万行
.fetchSize(2048) // 每次 fetch 2048 行
.parallelism(4) // 使用 4 个并行任务读取快照
.build();
效果: 大表全量同步时间从小时级降至分钟级。
注意:
parallelism仅影响快照阶段,binlog 读取仍为单并发。
9.2 快照阶段性能优化(snapshot.fetch.size, connection pool)
| 优化方向 | 配置项 | 说明 | 建议值 | 注意事项 |
|---|---|---|---|---|
| 批量读取 | scan.snapshot.fetch.size | 减少网络往返次数 | 2048 | 与数据库 max_allowed_packet 匹配 |
| 连接池 | jdbc.properties | 快照任务使用连接池 | 配置 HikariCP 或 Druid | Flink CDC 内部自动管理 |
| 并发控制 | scan.snapshot.parallelism | 显式设置快照并行度 | 等于 split 数量 | 避免过度并发压垮数据库 |
| 索引利用 | — | 分片字段需有索引 | 主键或唯一索引 | 否则分片扫描变全表扫描 |
| 内存配置 | TaskManager task.heap-size | 避免 OOM | 增大堆内存或启用堆外缓存 | 大宽表需更多内存 |
连接池配置示例(通过 JDBC 属性):
Map<String, String> jdbcProps = new HashMap<>();
jdbcProps.put("useSSL", "false");
jdbcProps.put("allowPublicKeyRetrieval", "true");
jdbcProps.put("autoDeserialize", "false");
// HikariCP 参数(Flink CDC 内部使用)
jdbcProps.put("dataSource.cachePrepStmts", "true");
jdbcProps.put("dataSource.prepStmtCacheSize", "250");
jdbcProps.put("dataSource.prepStmtCacheSqlLimit", "2048");
MySqlSource.<String>builder()
.jdbcProperties(jdbcProps)
// ... 其他配置
.build();
建议:
- 对 > 1000 万行的表,启用分片 + 并行快照。
- 监控数据库
Threads_connected、Select_scan指标,避免连接数爆炸。
9.3 减少对源库影响:read-only 用户、binlog 格式要求
| 措施 | 说明 | 配置/命令 | 注意事项 |
|---|---|---|---|
| 只读账号 | 避免误操作 | GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flink'@'%' | 不要授予 DROP、UPDATE 等权限 |
binlog_format | 必须为 ROW | SET GLOBAL binlog_format = 'ROW'; | STATEMENT/MIXED 无法捕获变更 |
binlog_row_image | 推荐 FULL | SET GLOBAL binlog_row_image = 'FULL'; | MINIMAL 可能丢失 before 值 |
read_only 模式 | 源库从节点开启 | SET GLOBAL read_only = ON; | 生产主库通常不开启 |
heartbeat_interval | 保持连接活跃 | heartbeat.interval.ms=30000 | 防止被防火墙断开 |
expire_logs_days | 保留足够 binlog | SET GLOBAL expire_logs_days = 7; | 至少保留 7 天,防止恢复失败 |
MySQL 推荐配置(my.cnf):
[mysqld]
server-id = 1
log-bin = mysql-bin
binlog-format = ROW
binlog-row-image = FULL
expire-logs-days = 7
read-only = 0 # 主库
安全原则:
- 使用专用账号,IP 白名单限制。
- 定期审计权限与连接日志。
- 避免在业务高峰期执行全量同步。
9.4 监控指标与常见瓶颈分析(Backpressure、Checkpoint Duration)
| 指标 | 位置 | 正常范围 | 异常表现 | 排查方法 |
|---|---|---|---|---|
| Backpressure | Flink Web UI(Task Metrics) | Green(无压力) | Yellow/Red | 检查下游 Sink 是否慢(如数据库写入瓶颈) |
| Checkpoint Duration | Checkpoint 面板 | < 间隔时间(如 5s) | > 30s | 检查状态大小、网络 I/O、存储性能 |
| Checkpoint Alignment | Subtask Details | 接近 0ms | > 1s | 表示数据倾斜或 Barrier 对齐等待 |
| Emit Delay | Source Metrics | < 100ms | > 1s | Source 读取慢,检查数据库负载 |
| Num Records in/out | Operator Metrics | 稳定波动 | 突增/突降 | 检查数据源或 Sink 是否异常 |
| State Size | Checkpoint Details | 稳定或缓慢增长 | 爆炸式增长 | 检查窗口未触发、状态未清理 |
常见瓶颈与解决方案:
| 瓶颈 | 现象 | 解决方案 |
|---|---|---|
| 数据库负载高 | CPU > 80%,慢查询增多 | 降低快照并行度、错峰同步、使用从库 |
| Sink 写入慢 | Backpressure 红色,延迟增长 | 增加 Sink 并行度、批量写入、优化索引 |
| Checkpoint 超时 | Failed Checkpoint | 增大超时时间、减少状态、优化网络 |
| OOM | TaskManager 重启 | 增大堆内存、减少 fetch.size、启用堆外状态 |
| 数据倾斜 | 某 subtask 处理数据远多于其他 | 优化 keyBy 字段、重新分区 |
监控建议:
- 集成 Prometheus + Grafana,建立实时监控看板。
- 设置告警规则:Checkpoint 失败 > 3 次、Backpressure 持续 5 分钟、延迟 > 1 分钟。
第10章 高级特性与扩展
10.1 Schema Evolution 支持现状与处理策略
| 变更类型 | Flink CDC 支持情况 | 处理策略 | 注意事项 |
|---|---|---|---|
| 新增字段 | 支持(after 包含新字段) | 下游需兼容,缺失字段设为 null | JSON 格式天然兼容 |
| 删除字段 | 支持(before 可能缺失) | 忽略或设为默认值 | 需注意反序列化器是否报错 |
| 字段重命名 | 不自动映射 | 需重建 Job 或使用视图 | 建议避免 |
| 类型变更 | 部分支持 | STRING → INT 可能失败 | 强类型语言(如 Java)易出错 |
| 主键变更 | 危险 | 可能导致分片策略失效 | 建议全量重建 |
| 表重命名 | 不感知 | 需重启 Job 并更新正则 | 无法动态发现 |
处理策略建议:
- 使用 JSON/Map 格式作为中间表示,避免强类型绑定。
- 在 Sink 层进行 schema 映射与转换(如使用 ProcessFunction)。
- 重大变更(如主键修改)建议:
- 停止 Job
- 清理状态(或从 Savepoint 恢复)
- 更新 Job 逻辑
- 重新部署
10.2 自定义 SourceSplitter 与 Reader
适用于需要自定义分片逻辑(如按时间分区、地理区域)的场景。
| 组件 | 作用 | 扩展方式 | 示例场景 |
|---|---|---|---|
| SourceSplitter | 将表划分为多个 SourceSplit | 实现 SourceSplitter<T> 接口 | 按 create_time 分月分片 |
| SplitEnumerator | 管理分片分配 | 实现 SplitEnumerator | 控制分片分发策略 |
| SourceReader | 读取单个分片数据 | 实现 SourceReader | 自定义数据解析逻辑 |
| SplitReader | 物理读取器 | 实现 SplitReader | 使用非 JDBC 方式读取(如 MyCat) |
扩展流程:
- 继承
MySqlSource或实现Source<OUT, S extends SourceSplit, ST>。 - 重写
createEnumerator()和createReader()。 - 注册自定义分片器与读取器。
注意:
- 属于高级开发,需深入 Flink Connector API。
- 建议优先使用现有配置,仅在必要时扩展。
10.3 Filter & Projection:在 Source 层过滤字段与记录
| 方式 | 说明 | 配置/代码 | 优点 |
|---|---|---|---|
| 字段过滤(Projection) | 只读取指定字段 | .columnTypes("id, name, age") | 减少网络传输与反序列化开销 |
| 记录过滤(Filter) | WHERE 条件过滤 | .startupOptions(StartupOptions.initial().withQuery("SELECT * FROM users WHERE status = 1")) | 减少数据量 |
| 正则过滤表 | 只同步匹配表 | .tableName("user_active|user_history") | 减少无关表同步 |
| Debezium filter | 使用 Debezium 内部过滤 | .debeziumProperties(Properties) / "column.exclude.list": "secret.*" | 支持列级过滤 |
代码示例:
// 只读取 id, name, email 字段
MySqlSource.<String>builder()
.columnTypes("id, name, email")
.build();
// 初始快照只读取 active 用户
StartupOptions options = StartupOptions.initial();
options.withQuery("SELECT id, name, status FROM users WHERE status = 'active'");
MySqlSource.<String>builder()
.startupOptions(options)
.build();
优势: 在源头减少数据量,提升整体性能。
安全: 可用于过滤敏感字段(如密码、身份证)。
10.4 支持 DDL 捕获(实验性功能)
| 功能 | 说明 | 启用方式 | 注意事项 |
|---|---|---|---|
| DDL 事件输出 | 将 ALTER TABLE 等操作作为事件输出 | .includeSchemaChanges(true) | 输出格式为 JSON,包含 ddl 字段 |
| DDL 事件结构 | 包含操作类型、SQL 语句、时间戳 | { "op": "c", "ddl": "ALTER TABLE ...", "ts_ms": ... } | op="c" 表示 schema change |
| 下游处理 | 需自定义逻辑处理 DDL | 使用 ProcessFunction 解析并执行 | 不支持自动应用到目标表 |
| 限制 | 仅捕获,不响应 | Flink Job 不会自动 reload schema | 需外部系统消费处理 |
| 典型用途 | 数据血缘、审计、元数据同步 | 写入 Kafka 或日志系统 | 用于监控表结构变更 |
启用 DDL 捕获:
MySqlSource.<String>builder()
.includeSchemaChanges(true) // 启用 DDL 事件
.deserializer(ChangelogJsonDeserializationSchema.INSTANCE)
.build();
DDL 事件示例:
{
"op": "c",
"ts_ms": 1712345678000,
"ddl": "ALTER TABLE users ADD COLUMN email VARCHAR(255) AFTER name",
"source": { "db": "test", "table": "users" }
}
警告: 该功能为实验性(experimental),API 可能变更,不建议在核心生产链路依赖。
建议用途: 辅助系统(如元数据中心、变更审计)。
第11章 常见问题与故障排查
11.1 权限不足导致连接失败
| 问题现象 | 根因分析 | 排查方法 | 解决方案 | 预防措施 |
|---|---|---|---|---|
Access denied for user 'flink'@'xxx' | 用户名/密码错误或 IP 未授权 | 查看 Flink 日志中的 SQLException | 核对用户名、密码、主机白名单 | 使用专用账号,配置 IP 白名单 |
SELECT command denied to user | 缺少 SELECT 权限 | 日志中出现 SELECT 查询失败 | 执行 GRANT SELECT ON db.table TO 'flink'@'%' | 最小权限原则,按需授权 |
RELOAD command denied | 无 RELOAD 权限,影响快照一致性 | 快照阶段报错 Cannot execute query | GRANT RELOAD ON *.* TO 'flink'@'%' | 若使用无锁快照可省略 |
REPLICATION SLAVE denied | 无 binlog 读取权限 | 无法启动 binlog reader,报 Access denied | GRANT REPLICATION SLAVE ON *.* TO 'flink'@'%' | 必须授权,否则无法增量同步 |
REPLICATION CLIENT denied | 无 SHOW MASTER STATUS 权限 | 无法获取当前 binlog 位点 | GRANT REPLICATION CLIENT ON *.* TO 'flink'@'%' | 必须授权 |
SHOW DATABASES denied | 无数据库列表权限 | 启动时报 SHOW DATABASES 错误 | GRANT SHOW DATABASES ON *.* TO 'flink'@'%' | 若指定 databaseName 可省略 |
推荐最小权限集合:
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flink'@'%' IDENTIFIED BY 'cdc_password';
FLUSH PRIVILEGES;
验证权限命令:
mysql -h your-host -u flink -p -e "SHOW MASTER STATUS; SELECT 1 FROM your_table LIMIT 1;"
11.2 Binlog 不可用或格式不兼容
| 问题现象 | 根因分析 | 排查方法 | 解决方案 | 预防措施 |
|---|---|---|---|---|
binlog format statement is not supported | binlog_format 不是 ROW | 日志中明确提示 | 登录 MySQL 执行:SET GLOBAL binlog_format = 'ROW'; | 部署前检查配置 |
binlog_row_image is minimal | before 字段缺失 | CDC 数据中 before 为 null | SET GLOBAL binlog_row_image = 'FULL'; | 生产环境必须设为 FULL |
Could not find first log file name | binlog 文件被清理 | 启动时报 Could not find first log file | 重新全量同步 / 从备份恢复 binlog | 设置 expire_logs_days = 7 |
Could not read from offset | 位点超出保留范围 | Checkpoint 中记录的位点已不存在 | 清除状态重新同步或使用 Savepoint 回退 | 启用外部化 Checkpoint,定期备份 |
Unknown binlog event type: FORMAT_DESCRIPTION | 版本兼容问题(极少见) | 特定 MySQL 版本与 Debezium 不兼容 | 升级 Flink CDC 或 MySQL | 使用主流版本(如 MySQL 5.7/8.0) |
MySQL Binlog 检查命令:
-- 检查格式
SHOW VARIABLES LIKE 'binlog_format';
SHOW VARIABLES LIKE 'binlog_row_image';
-- 查看当前 binlog 文件
SHOW MASTER STATUS;
-- 查看可用 binlog
SHOW BINARY LOGS;
-- 设置保留7天
SET GLOBAL expire_logs_days = 7;
注意: 修改
binlog_format需重启连接,已建立的连接不生效。
11.3 快照锁表问题与无锁快照配置
| 问题现象 | 根因分析 | 排查方法 | 解决方案 | 配置说明 |
|---|---|---|---|---|
数据库响应变慢,SHOW PROCESSLIST 显示 Waiting for table metadata lock | Flink CDC 使用 FLUSH TABLES WITH READ LOCK 获取一致性快照 | 监控数据库锁等待 | 启用无锁快照(Lock-free Snapshot) | 默认开启(Flink CDC 2.3+) |
FLUSH TABLES WITH READ LOCK 失败 | 用户无 RELOAD 权限或主库不允许 | 日志中出现 FLUSH 命令拒绝 | 禁用全局锁,使用无锁快照 | .scan.incremental.snapshot.enabled(true) |
| 快照期间写入阻塞 | 使用了全局读锁 | 业务 INSERT/UPDATE 延迟升高 | 确保使用 READ UNCOMMITTED 隔离级别 | 无锁快照基于 MVCC 实现 |
| 大表快照耗时过长 | 单线程全表扫描 | TaskManager 日志显示长时间读取 | 启用分片 + 并行快照 | .splitSize(10000).parallelism(4) |
无锁快照原理:
- 基于 RR(Repeatable Read)隔离级别和主键范围扫描。
- 按主键分片,逐个读取,不阻塞写入。
- 保证每个分片内部一致性,但非全局瞬时一致性。
启用无锁快照配置:
MySqlSource.<String>builder()
.databaseName("prod")
.tableName("large_table")
.scanStartupMode(StartupMode.INITIAL)
// 启用增量快照(即无锁快照)
.scanIncrementalSnapshotEnabled(true)
// 分片配置
.splitSize(20000)
.parallelism(4)
.build();
优势: 对源库影响极小,适合生产环境。
注意: 若业务要求跨表事务一致性,需额外处理。
11.4 数据重复或丢失的根因分析
| 问题类型 | 现象 | 根因 | 排查方法 | 解决方案 |
|---|---|---|---|---|
| 数据重复 | 目标表出现多条相同主键记录 | 1. Sink 不支持幂等 2. Checkpoint 失败后回退 | 检查 Sink 写入逻辑、Checkpoint 失败日志 | 使用事务型 Sink(如 Kafka 事务、Doris) |
| 数据丢失 | 源库有变更,目标无反映 | 1. Binlog 被清理 2. Source 过滤条件过严 3. 反序列化失败静默丢弃 | 检查 binlog 保留、日志中是否有解析错误 | 启用外部化 Checkpoint,配置日志告警 |
| 更新丢失 | UPDATE 操作未生效 | 1. Sink 未实现 upsert 2. 主键不唯一导致插入而非更新 | 检查目标表主键、Sink SQL 逻辑 | 使用 ON DUPLICATE KEY UPDATE 或 MergeTree |
| Delete 未同步 | DELETE 操作未传播 | 1. Sink 忽略 delete 事件 2. 目标表不支持 delete | 检查 Sink 是否处理 op=d | 确保 Sink 支持全量变更类型 |
| 全量阶段重复 | 全量 + 增量衔接处数据重复 | 快照结束与 binlog 起始位点重叠 | 检查位点衔接逻辑 | Flink CDC 自动去重(基于 source.ts_ms) |
详细根因与对策:
| 场景 | 根因 | 解决方案 |
|---|---|---|
| Checkpoint 失败导致重复 | 任务失败前已提交 Sink 但未完成 Checkpoint | 使用两阶段提交(2PC) Sink(如 Kafka) |
| 反序列化异常静默丢弃 | JSON 解析失败,默认行为是跳过 | 自定义反序列化器,记录错误日志或发到侧输出流 |
| 并行度 > 1 导致乱序 | 多个 Source 任务读取 binlog | 禁止设置 MySqlSource 并行度 > 1 |
| Savepoint 恢复时从头开始 | 未正确指定启动模式 | 使用 StartupMode.SPECIFIC_OFFSETS 指定位点 |
| 网络抖动导致重试 | Source 重试机制 | 确保 Sink 幂等或事务性 |
推荐保障手段:
- 端到端一致性:启用 Checkpoint + 事务型 Sink。
- 监控告警:监控
numRecordsInvsnumRecordsOut、Checkpoint 失败次数。 - 数据校验:定期对账(如源库 COUNT 与目标库 COUNT 对比)。