一、问题前置:你遇到的“假正常”陷阱

做数据同步的人大概率踩过这个坑: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端的位点在推进,但下游还是没数据,那要检查是不是配置了过滤规则,把所有数据都滤掉了。比如你只同步表userid>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都是500800,那过滤后自然没数据流到下游。

三、第二阶段:检查中间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代码,如果插入的字段类型不匹配(比如上游的idint,下游的表的idvarchar),插入会失败,事务回滚,下游自然收不到数据。

踩坑例子:有个开发同步的上游表的idint,下游的表的idvarchar,插入时类型不匹配,提交失败,他的代码里没打印异常日志,结果所有数据都没提交到下游。

五、应用场景、优缺点、注意事项

5.1 应用场景

这套排查方法适用于所有Flink CDC的同步场景,比如同步MySQL到Kafka、同步MySQL到MySQL、同步MySQL到ClickHouse等,只要是Flink CDC任务假正常、下游收不到数据的场景,都可以用这套方法排查。

5.2 优缺点

优点:排查路径清晰,从左到右顺着数据流动的顺序,每一步都有具体的操作和例子,容易落地;覆盖了Source、Channel、Sink三个核心环节,不会漏掉关键的问题点。 缺点:需要对Flink CDC的整个链路有一定的了解,对于完全没接触过Flink的人来说,可能需要先补一些基础;排查过程需要登录源端数据库、查看Flink的WebUI和日志,有一定的操作成本。

5.3 注意事项

  1. 排查时要按顺序来:先查Source,再查Channel,最后查Sink,不要跳着查,否则可能漏掉关键的问题点。
  2. 要关注日志的细节:很多问题的线索都在日志里,比如Source的位点、序列化的异常、Sink的提交日志等,排查时要仔细看日志。
  3. 要配置好异常处理:在Source的过滤、序列化、Sink的提交等环节,都要配置好异常处理,打印出异常日志,这样排查问题时能更快定位。

六、文章总结

Flink CDC任务假正常、下游收不到数据的问题,本质是数据流动的某个环节出了问题,只是这个环节的问题没有触发明显的报错。只要顺着数据从Source到Sink的完整链路,一段一段抠细节,就能找到问题所在。

排查时要注意,不要只看Flink的WebUI的Running状态,要深入到每个环节的具体状态:Source端要检查是否真的连上Binlog、位点是否推进、过滤规则是否误杀;中间Channel要检查序列化是否出错、缓冲区是否满了;Sink端要检查事务是否提交、提交是否失败。

只要掌握了这套排查方法,遇到这类问题就能快速定位和解决,不会再陷入“任务看起来正常,就是没数据”的困境。