一、问题前置:你遇到的“假正常”陷阱
做数据同步的人大概率踩过这个坑:Flink CDC任务跑在那,监控面板上CPU、内存、线程数全是绿的,日志里也没ERROR没WARN,可下游的数据库、消息队列就是死活收不到新数据——明明源端的表一直在增删改,同步任务却像“假死”了一样。
这种问题比直接报错难查多了:报错能定位到具体行,可假正常的问题得顺着数据流动的路径,一段一段抠细节。今天就按数据从源端到下游的完整链路,给你一套可落地的排查步骤,每一步都带真实场景的例子。
二、排查链路:从Source到Sink的三段式检查
先明确整个数据流动的核心链路:源端数据库(比如MySQL)的Binlog → Flink CDC的Source算子(负责拉取Binlog) → 中间Channel(负责数据传输、序列化) → Sink算子(负责把数据写入下游)。我们就按这个顺序,从左到右拆成三个阶段查。
2.1 第一阶段:检查Source端的“真活跃”
很多人看Source活跃,只看Flink的WebUI里Source算子的Running状态,这完全不够——有些场景下Source算子是活的,但根本没在拉取新的Binlog。
2.1.1 核心检查点1:Binlog是否真的被读取
Flink CDC的Source算子,本质是模拟一个MySQL的Slave(从库),从主库拉取Binlog。如果Source真的在工作,那它会不断向主库发送心跳包,主库的Slave状态里一定会有这个Flink任务的连接。
操作步骤(拿MySQL举例): 登录源端MySQL,执行查看所有从库连接的命令:
-- 查看当前主库的所有从库连接状态
SHOW SLAVE HOSTS;
如果这个命令的输出里,没有Flink任务的IP(或者你配置的Source端服务器IP),那说明Source算子根本没连上主库的Binlog,或者连了但已经断了。
举个真实踩坑的例子:之前有个开发把Flink CDC的Source配置里的server-id设成了和主库现有从库一样的值(MySQL要求每个Slave的server-id必须唯一),结果主库会把重复server-id的连接踢掉,Source算子看似是Running状态,实则已经丢了Binlog连接。
2.1.2 核心检查点2:Source的位点是否在推进
就算Source连上了,也可能因为位点卡壳,拉不到新数据。比如源端的Binlog已经滚到了下一个文件,但Source还停在旧文件的末尾,或者旧文件已经被主库清理了。
操作步骤: 先看Flink CDC的Source端配置的位点(如果是从指定位点启动的话),再对比主库当前的Binlog状态:
-- 查看主库当前正在写入的Binlog文件和位置
SHOW MASTER STATUS;
假设主库的输出是:当前Binlog文件是mysql-bin.000005,位置是12345。那再看Flink CDC任务的日志,找Source端的位点信息——Flink CDC启动时会打印类似这样的日志:
[INFO] 正在读取Binlog文件: mysql-bin.000004, 位置: 99999
如果Source的位点还停在mysql-bin.000004的末尾,说明它卡在了旧文件,没跟着主库的Binlog滚动。
踩坑例子:有个开发把Flink CDC的scan.startup.mode设成了initial(初始化模式),但源端表的数据量特别大,初始化扫描花了12个小时,这期间主库的Binlog已经滚了5个文件,Source算子的位点一直停在初始化扫描的位置,初始化完成后才发现,之前的新数据对应的Binlog已经被主库清理了(主库的Binlog过期时间设的是8小时),结果初始化后就再也拉不到新数据了。
2.1.3 核心检查点3:过滤规则是否误杀
如果Source端的位点在推进,但下游还是没数据,那要检查是不是配置了过滤规则,把所有数据都滤掉了。比如你只同步表user的id>1000的数据,但源端的user表最近插入的都是id<1000的数据,自然就没数据流到下游。
操作步骤: 找到Flink CDC Source端的过滤配置,比如用Java写的Flink CDC作业,代码里的过滤逻辑是这样的:
// 仅同步id>1000的用户数据
SourceFunction<String> sourceFunction = MySQLSource.<String>builder()
.hostname("localhost")
.port(3306)
.username("root")
.password("123456")
.databaseList("test_db")
.tableList("test_db.user")
.deserializer(new JsonDebeziumDeserializationSchema())
.startupOptions(StartupOptions.latest())
.build();
// 过滤逻辑
DataStream<String> stream = env.addSource(sourceFunction)
.filter(json -> {
try {
// 解析JSON,获取id的值
ObjectNode node = (ObjectNode) new ObjectMapper().readTree(json);
int id = node.get("after").get("id").asInt();
// 仅保留id>1000的数据
return id > 1000;
} catch (Exception e) {
return false;
}
});
如果源端最近插入的user数据的id都是500、800,那过滤后自然没数据流到下游。
三、第二阶段:检查中间Channel的传输和序列化
就算Source端的位点在推进,也有数据产生,中间的Channel环节也可能出问题——比如数据在序列化时出错,或者Channel的缓冲区满了,导致数据传不下去。
3.1 核心检查点1:序列化是否出错
Flink的Source算子产生的数据,要先序列化成二进制,再通过Channel传给下游的Sink算子。如果序列化时出错,Flink会默认把这条数据丢掉,而且不会打印ERROR日志(除非你配置了更严格的异常处理)。
操作步骤: 找到Flink作业的序列化配置,比如用JSON序列化的场景,代码是这样的:
// JSON序列化器
public class JsonDebeziumDeserializationSchema implements DebeziumDeserializationSchema<String> {
private static final ObjectMapper mapper = new ObjectMapper();
@Override
public void deserialize(SourceRecord record, Collector<String> out) throws Exception {
try {
// 把Debezium的SourceRecord转成JSON字符串
String json = mapper.writeValueAsString(record.value());
out.collect(json);
} catch (Exception e) {
// 这里如果不打印日志,序列化出错的话,数据会被默默丢掉
// e.printStackTrace();
}
}
@Override
public TypeInformation<String> getProducedType() {
return BasicTypeInfo.STRING_TYPE_INFO;
}
}
如果源端的表有特殊数据类型(比如MySQL的geometry类型),Debezium转成的SourceRecord里的这个字段,JSON序列化时会出错,结果这条数据就被丢掉了。
踩坑例子:有个开发同步的表有geometry类型的字段,序列化时出错,他的代码里没打印异常日志,结果所有带这个字段的Binlog都被丢掉了,下游自然收不到数据。
3.2 核心检查点2:Channel缓冲区是否满了
Flink的Channel有个缓冲区大小的配置,如果上游Source产生数据的速度远大于下游Sink处理的速度,缓冲区会被占满,Source算子会暂停产生新数据,导致下游收不到。
操作步骤: 先看Flink作业的WebUI,找到Source算子和Sink算子之间的Channel,看缓冲区的使用率。如果使用率一直是100%,说明缓冲区满了。
再看Flink的配置,缓冲区的默认大小是32MB,如果你的数据量特别大,需要调大这个配置:
# Flink配置文件里的缓冲区大小配置
taskmanager.network.memory.fraction=0.2
taskmanager.network.memory.min=64MB
taskmanager.network.memory.max=256MB
踩坑例子:有个开发同步的表每条数据都有大文本字段,单条数据大小是10MB,Flink的缓冲区默认是32MB,只能存3条数据,而下游Sink处理一条数据需要2秒,结果缓冲区很快满了,Source算子暂停产生数据,下游就收不到新数据了。
四、第三阶段:检查Sink端的事务提交
就算中间Channel把数据传给了Sink算子,Sink端的事务提交也可能出问题——比如事务没提交,或者提交失败,导致下游收不到数据。
4.1 核心检查点1:事务是否提交
Flink的Sink算子,比如同步到MySQL的Sink,默认是批量提交事务的——只有当批量的数据量达到阈值,或者时间达到阈值,才会提交事务。如果批量的阈值设得特别大,或者时间阈值设得特别长,那Sink算子会攒着数据不提交,下游自然收不到。
操作步骤: 找到Sink端的批量配置,比如同步到MySQL的Sink,代码是这样的:
// 批量提交的配置
JdbcSink.sink(
"INSERT INTO user (id, name) VALUES (?, ?)",
(statement, row) -> {
statement.setInt(1, row.getField(0));
statement.setString(2, row.getField(1));
},
JdbcExecutionOptions.builder()
// 批量提交的条数:攒1000条才提交
.withBatchSize(1000)
// 批量提交的时间间隔:攒满1000条或者过了60秒才提交
.withBatchIntervalMs(60000)
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:mysql://localhost:3306/test_db")
.withDriverName("com.mysql.cj.jdbc.Driver")
.withUsername("root")
.withPassword("123456")
.build()
);
如果源端的表最近插入的频率特别低,比如每5秒才插一条,那Sink算子要攒满1000条需要5000秒(大概1.4小时),或者过了60秒才提交,结果下游要等很久才能收到数据。
4.2 核心检查点2:事务提交是否失败
就算事务达到了提交条件,也可能提交失败——比如下游的表结构和上游的不一样,导致插入失败,Flink会回滚事务,而且不会打印ERROR日志(除非你配置了异常处理)。
操作步骤:
找到Sink端的异常处理配置,比如同步到MySQL的Sink,如果插入失败,代码里没打印异常,那提交失败的日志就不会显示。比如上面的JdbcSink代码,如果插入的字段类型不匹配(比如上游的id是int,下游的表的id是varchar),插入会失败,事务回滚,下游自然收不到数据。
踩坑例子:有个开发同步的上游表的id是int,下游的表的id是varchar,插入时类型不匹配,提交失败,他的代码里没打印异常日志,结果所有数据都没提交到下游。
五、应用场景、优缺点、注意事项
5.1 应用场景
这套排查方法适用于所有Flink CDC的同步场景,比如同步MySQL到Kafka、同步MySQL到MySQL、同步MySQL到ClickHouse等,只要是Flink CDC任务假正常、下游收不到数据的场景,都可以用这套方法排查。
5.2 优缺点
优点:排查路径清晰,从左到右顺着数据流动的顺序,每一步都有具体的操作和例子,容易落地;覆盖了Source、Channel、Sink三个核心环节,不会漏掉关键的问题点。 缺点:需要对Flink CDC的整个链路有一定的了解,对于完全没接触过Flink的人来说,可能需要先补一些基础;排查过程需要登录源端数据库、查看Flink的WebUI和日志,有一定的操作成本。
5.3 注意事项
- 排查时要按顺序来:先查Source,再查Channel,最后查Sink,不要跳着查,否则可能漏掉关键的问题点。
- 要关注日志的细节:很多问题的线索都在日志里,比如Source的位点、序列化的异常、Sink的提交日志等,排查时要仔细看日志。
- 要配置好异常处理:在Source的过滤、序列化、Sink的提交等环节,都要配置好异常处理,打印出异常日志,这样排查问题时能更快定位。
六、文章总结
Flink CDC任务假正常、下游收不到数据的问题,本质是数据流动的某个环节出了问题,只是这个环节的问题没有触发明显的报错。只要顺着数据从Source到Sink的完整链路,一段一段抠细节,就能找到问题所在。
排查时要注意,不要只看Flink的WebUI的Running状态,要深入到每个环节的具体状态:Source端要检查是否真的连上Binlog、位点是否推进、过滤规则是否误杀;中间Channel要检查序列化是否出错、缓冲区是否满了;Sink端要检查事务是否提交、提交是否失败。
只要掌握了这套排查方法,遇到这类问题就能快速定位和解决,不会再陷入“任务看起来正常,就是没数据”的困境。
评论
围绕“Flink CDC任务看似运行正常但下游始终收不到Binlog变更,沿着Source活跃状态、Channel序列化到Sink事务提交逐层检查位点推进并输出完整排查路径”参与讨论