Article

日志采集 Flink CDC 插件

更新于:2026-07-13

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 流式读取。确保历史数据不丢失,实现全量 + 增量一体化同步。大表快照可能影响源库性能,建议在低峰期执行。
概念名称说明用途注意事项
全量 + 增量一体化支持首次全量读取表数据,自动切换到增量 binlog 模式,无需手动切换。简化数据同步流程,避免数据断层。需保证 binlog 保留时间足够长,覆盖快照完成时间。
无锁读取(Lock-free Snapshot)使用 RR 隔离级别和 SELECT ... FOR SHAREREAD UNCOMMITTED 实现快照,避免锁表。减少对生产库的影响,提升并发能力。需数据库支持 MVCC(如 InnoDB)。
端到端 Exactly-Once 语义结合 Flink Checkpoint 机制,确保每条变更仅被处理一次。保障数据一致性,适用于金融、订单等关键业务。Sink 端需支持幂等写入或事务。
Schema Evolution 支持支持表结构变更(如加列)的自动感知(部分支持)。适应业务迭代,减少运维成本。当前对 DDL 的完整支持仍在演进中。
应用场景:实时数仓将业务库数据实时同步到数仓(如 Doris、ClickHouse、Iceberg)。构建低延迟的数据分析平台。需注意目标系统写入性能。
应用场景:缓存更新捕获数据库变更后,自动更新 Redis 缓存。避免缓存与数据库不一致。需设计合理的缓存失效策略。
应用场景:微服务解耦将数据库变更作为事件发布,驱动其他服务响应。实现事件驱动架构(EDA)。需结合消息中间件(如 Kafka)进行解耦。
对比项Flink CDC传统方案(如 DataX、Sqoop)说明注意事项
同步模式实时流式同步(秒级延迟)批处理同步(分钟/小时级延迟)Flink CDC 基于流处理,延迟更低。传统方案适合离线批处理场景。
架构复杂度简单:Flink Job 直接读取数据库日志复杂:常需 Kafka、Kafka Connect 中转Flink CDC 内嵌 Debezium,无需额外组件。减少运维成本。
数据一致性支持端到端 Exactly-Once通常为 At-Least-Once,易重复Flink Checkpoint 保障一致性。Sink 需支持事务或幂等。
全量增量一体化支持自动切换需分别配置全量和增量任务简化任务管理。避免数据重复或遗漏。
资源占用持续运行,占用一定资源按需运行,资源占用低适合长期运行的实时任务。需合理配置资源。
扩展性支持并行读取、多表同步扩展性差,通常单线程可水平扩展应对大数据量。需合理设置分片参数。
数据源支持状态所需依赖关键特性注意事项
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)

工具版本要求安装方式用途注意事项
JavaJDK 8 或 JDK 11官网下载并配置 JAVA_HOMEFlink 运行基础环境推荐使用 OpenJDK,避免商业授权问题
Maven3.5+官网下载并配置 PATH项目依赖管理与构建配置国内镜像(如阿里云)提升下载速度
IDEIntelliJ IDEA / Eclipse官网下载代码编写与调试推荐 IDEA,对 Flink 有良好支持
Flink1.13+(推荐 1.17+)官网下载或通过 Maven 依赖流处理引擎Flink CDC 3.x 要求 Flink 1.17+
构建命令mvn clean package命令行执行打包项目为 jar确保 pom.xml 正确配置依赖
运行模式配置方式适用场景启动命令注意事项
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-applicationJar 包需包含依赖(fat jar)
Session 模式先启动集群,再提交任务多任务共享集群flink run -d job.jar任务间可能资源竞争
依赖名称用途注意事项
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(生产部署时由集群提供)。

方法/配置语法用途代码示例注意事项
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
执行 Jobenv.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(时间戳)、beforeafter 等字段。

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_recoveryinitial=全量+增量,schema_only=仅结构
snapshot.locking.mode.debeziumConfig("snapshot.locking.mode", "none")快照锁策略none, minimal, extendednone 表示无锁读取(推荐)
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(保留精度), stringstring 更安全
skip.messages.debeziumConfig("skip.messages", "true")跳过无法解析的消息true / false避免作业失败
heartbeat.interval.ms.debeziumConfig("heartbeat.interval.ms", "30000")心跳发送间隔30000(30秒)用于监控连接状态

调优建议:

  • 生产环境建议设置 snapshot.locking.mode=none 避免锁表。
  • 大表同步可调大 binlog.buffer.sizepoll.interval.ms 降低 CPU。
  • 使用 decimal.handling.mode=string 防止精度丢失。

3.4 启动模式(StartupMode)详解:initial、latest-offset、timestamp、specific-offset

启动模式语法用途代码示例注意事项
INITIALStartupMode.INITIAL首次启动:先全量快照,再读增量 binlog.startupMode(StartupMode.INITIAL)适用于新表同步
LATEST_OFFSETStartupMode.LATEST_OFFSET仅从当前最新位点开始读取 binlog.startupMode(StartupMode.LATEST_OFFSET)不读历史数据,常用于灾备恢复
TIMESTAMPStartupMode.TIMESTAMP从指定时间戳开始读取.startupMode(StartupMode.TIMESTAMP) / .startupTimestamp(1672531200L)时间戳单位为秒
SPECIFIC_OFFSETStartupMode.SPECIFIC_OFFSET从指定 binlog 文件和位置开始.startupMode(StartupMode.SPECIFIC_OFFSET) / .startupSpecificOffset("mysql-bin.000003", 154L)精准恢复场景
NEVERStartupMode.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 JSONChangelogJsonDeserializationSchemaFlink 自定义格式,兼容 changelog stream["+I",{"id":1,"name":"Charlie"}]"+I"=insert, "-U"=旧值 update, "+U"=新值 update, "-D"=delete
自定义 RowDataRowDataDeserializationSchema输出为 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
  }
}

注意事项:

  • beforeafter 是嵌套对象,需递归解析。
  • 时间字段可能为字符串或时间戳,取决于配置。
  • 大字段(如 BLOB)可能被截断或 Base64 编码。

4.2 如何解析 insert/update/delete 事件

事件类型判断条件处理逻辑代码示例(Java)注意事项
Insertop == "c""+I"取 after 字段作为新数据if (op.equals("c")) { User user = parseAfter(json); }快照中的 insert 也标记为 "c"
Updateop == "u"before 为旧值,after 为新值if (op.equals("u")) { User old = parseBefore(json); User updated = parseAfter(json); }区分旧值和新值
Deleteop == "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 输出的 JSONJsonObject 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();

注意事项:

  • 水印生成必须在 keyBywindow 之前调用。
  • 延迟时间(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 流程:

  1. Flink 触发 Checkpoint。
  2. Source 将当前 binlog 位点作为状态保存。
  3. Barrier 流经所有算子,触发状态快照。
  4. Sink 执行预提交(pre-commit),保存事务状态。
  5. JobManager 确认所有任务快照成功。
  6. 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 测试

恢复流程:

  1. 作业重启。
  2. Flink 从持久化存储加载最近的 Checkpoint/Savepoint。
  3. MySqlSource 读取其中的 binlog 位点(offset)。
  4. 连接 MySQL,从该位点继续读取 binlog。
  5. 继续处理后续变更。

注意事项:

  • 必须确保 MySQL 的 expire_logs_daysbinlog_expire_logs_seconds 足够长,覆盖最大恢复时间。
  • 若 binlog 已被清理,将导致恢复失败,需重新全量同步。

6.4 并行度与分片策略对一致性的影响

策略说明一致性影响注意事项
MySqlSource 并行度=1单任务读取 binlog强一致性,事件顺序完全保序吞吐受限,适用于中小数据量
并行度>1Flink 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.versionDebezium 结构版本"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 CDCConnect 支持动态发现新表真正动态 / 脱离 Flink 生态适合已使用 Kafka 的架构
Flink Application Mode + Savepoint修改 Job 逻辑并从 Savepoint 恢复flink run -s <savepoint> ...保证状态恢复 / 需重新编译打包适合版本升级

建议: 优先采用”多 Job + 配置化”方案,平衡灵活性与稳定性。

7.4 数据路由:将不同表写入不同 Sink

方法说明代码示例注意事项
sideOutput + ProcessFunction使用侧输出流分流见下方完整示例最灵活,推荐使用
split() (Deprecated)已废弃,不推荐避免使用
多 DataStream 分支map/filter 后分别 addSinkusersStream.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 ProducerFlinkKafkaProducer<String>需引入 flink-connector-kafka
序列化格式建议使用 JSON 或 AVROChangelog JSON、Debezium JSON便于下游消费
Key 设计使用主键哈希提升并行度.setKeyedSerializationSchema(...)避免热点
事务写入保障 Exactly-Once.setWriteTransactionsTimeout(...)需 Kafka 0.11+
Topic 路由按表名映射到不同 Topicdb.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 类型JdbcSinkJdbcSink.sink(...)支持批处理写入
批量大小控制写入频率.withBatchSize(1000)平衡延迟与吞吐
批量间隔超时强制提交.withBatchIntervalMs(2000)防止数据积压
更新策略insert or updateON 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:构建湖仓一体架构

特性HudiIceberg说明
写入模式COPY_ON_WRITE / MERGE_ON_READAPPEND / UPSERT支持 CDC 增量写入
Flink 集成HoodieFlinkStreamerFlinkIcebergSink均支持流式写入
更新支持支持 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 分析

对比项DorisClickHouse
写入方式Stream Load(HTTP)HTTP Interface / Native TCP
Flink ConnectorDorisSinkClickHouseSink(社区)
UPSERT 支持唯一模型表CollapsingMergeTree / ReplacingMergeTree
实时性秒级秒级
语法兼容MySQLSQL-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.size1024每次从数据库 fetch 的行数增大至 512~4096 提升吞吐受 JDBC fetchSize 限制,过大可能占用内存
scan.split.size8096每个分片(split)的行数大表设为 10000~50000,提升并行度值越小,并行任务越多,但小文件多
scan.incremental.snapshot.chunk.size1024增量快照分块大小(行)控制 binlog 回放粒度split.size 类似,用于增量阶段
parallelism1(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 或 DruidFlink 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_connectedSelect_scan 指标,避免连接数爆炸。

9.3 减少对源库影响:read-only 用户、binlog 格式要求

措施说明配置/命令注意事项
只读账号避免误操作GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flink'@'%'不要授予 DROP、UPDATE 等权限
binlog_format必须为 ROWSET GLOBAL binlog_format = 'ROW';STATEMENT/MIXED 无法捕获变更
binlog_row_image推荐 FULLSET GLOBAL binlog_row_image = 'FULL';MINIMAL 可能丢失 before 值
read_only 模式源库从节点开启SET GLOBAL read_only = ON;生产主库通常不开启
heartbeat_interval保持连接活跃heartbeat.interval.ms=30000防止被防火墙断开
expire_logs_days保留足够 binlogSET 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)

指标位置正常范围异常表现排查方法
BackpressureFlink Web UI(Task Metrics)Green(无压力)Yellow/Red检查下游 Sink 是否慢(如数据库写入瓶颈)
Checkpoint DurationCheckpoint 面板< 间隔时间(如 5s)> 30s检查状态大小、网络 I/O、存储性能
Checkpoint AlignmentSubtask Details接近 0ms> 1s表示数据倾斜或 Barrier 对齐等待
Emit DelaySource Metrics< 100ms> 1sSource 读取慢,检查数据库负载
Num Records in/outOperator Metrics稳定波动突增/突降检查数据源或 Sink 是否异常
State SizeCheckpoint Details稳定或缓慢增长爆炸式增长检查窗口未触发、状态未清理

常见瓶颈与解决方案:

瓶颈现象解决方案
数据库负载高CPU > 80%,慢查询增多降低快照并行度、错峰同步、使用从库
Sink 写入慢Backpressure 红色,延迟增长增加 Sink 并行度、批量写入、优化索引
Checkpoint 超时Failed Checkpoint增大超时时间、减少状态、优化网络
OOMTaskManager 重启增大堆内存、减少 fetch.size、启用堆外状态
数据倾斜某 subtask 处理数据远多于其他优化 keyBy 字段、重新分区

监控建议:

  • 集成 Prometheus + Grafana,建立实时监控看板。
  • 设置告警规则:Checkpoint 失败 > 3 次、Backpressure 持续 5 分钟、延迟 > 1 分钟。

第10章 高级特性与扩展

10.1 Schema Evolution 支持现状与处理策略

变更类型Flink CDC 支持情况处理策略注意事项
新增字段支持(after 包含新字段)下游需兼容,缺失字段设为 nullJSON 格式天然兼容
删除字段支持(before 可能缺失)忽略或设为默认值需注意反序列化器是否报错
字段重命名不自动映射需重建 Job 或使用视图建议避免
类型变更部分支持STRING → INT 可能失败强类型语言(如 Java)易出错
主键变更危险可能导致分片策略失效建议全量重建
表重命名不感知需重启 Job 并更新正则无法动态发现

处理策略建议:

  • 使用 JSON/Map 格式作为中间表示,避免强类型绑定。
  • 在 Sink 层进行 schema 映射与转换(如使用 ProcessFunction)。
  • 重大变更(如主键修改)建议:
    1. 停止 Job
    2. 清理状态(或从 Savepoint 恢复)
    3. 更新 Job 逻辑
    4. 重新部署

10.2 自定义 SourceSplitter 与 Reader

适用于需要自定义分片逻辑(如按时间分区、地理区域)的场景。

组件作用扩展方式示例场景
SourceSplitter将表划分为多个 SourceSplit实现 SourceSplitter<T> 接口create_time 分月分片
SplitEnumerator管理分片分配实现 SplitEnumerator控制分片分发策略
SourceReader读取单个分片数据实现 SourceReader自定义数据解析逻辑
SplitReader物理读取器实现 SplitReader使用非 JDBC 方式读取(如 MyCat)

扩展流程:

  1. 继承 MySqlSource 或实现 Source<OUT, S extends SourceSplit, ST>
  2. 重写 createEnumerator()createReader()
  3. 注册自定义分片器与读取器。

注意:

  • 属于高级开发,需深入 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 queryGRANT RELOAD ON *.* TO 'flink'@'%'若使用无锁快照可省略
REPLICATION SLAVE denied无 binlog 读取权限无法启动 binlog reader,报 Access deniedGRANT REPLICATION SLAVE ON *.* TO 'flink'@'%'必须授权,否则无法增量同步
REPLICATION CLIENT deniedSHOW 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 supportedbinlog_format 不是 ROW日志中明确提示登录 MySQL 执行:SET GLOBAL binlog_format = 'ROW';部署前检查配置
binlog_row_image is minimalbefore 字段缺失CDC 数据中 before 为 nullSET GLOBAL binlog_row_image = 'FULL';生产环境必须设为 FULL
Could not find first log file namebinlog 文件被清理启动时报 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 lockFlink 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。
  • 监控告警:监控 numRecordsIn vs numRecordsOut、Checkpoint 失败次数。
  • 数据校验:定期对账(如源库 COUNT 与目标库 COUNT 对比)。