You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Canal同步MySQL数据到Flink时遇NOT NULL字段写入空值问题求助

问题分析与解决方案

问题根源

你的Kafka消息是Canal输出的嵌套格式(业务数据包裹在data数组内),但Flink定义的transactions表直接在JSON顶层查找account_id等字段,导致解析失败,这些字段值变为null。而MySQL的spend_report表中account_id是NOT NULL约束,写入时触发异常。


方案一:修正Kafka表的JSON解析逻辑(适配Canal格式)

修改transactions表定义,通过JSON路径提取data数组中的第一个元素作为业务数据行:

tEnv.executeSql(
        "create table transactions \n" +
                "(\n" +
                "account_id bigint,\n" +
                "amount bigint,\n" +
                "transaction_time timestamp(3),\n" +
                "watermark for transaction_time as transaction_time - interval '5' second\n" +
                ") with \n" +
                "(\n" +
                "'connector'='kafka',\n" +
                "'topic'='transactions',\n" +
                "'properties.bootstrap.servers'='hadoop102:9092',\n" +
                "'scan.startup.mode'='latest-offset',\n" +
                "'format'='json',\n" +
                "'json.fail-on-missing-field' = 'false',\n" +
                "'json.ignore-parse-errors' = 'true',\n" +
                // 指定JSON根路径为data数组的第一个元素
                "'json.path'='$.data[0]'\n" +
                ");"
);

此配置让Flink从data数组内读取业务字段,避免解析出null值。


使用专门适配Canal格式的连接器,自动解析嵌套结构,无需手动处理JSON路径:

  1. 添加Maven依赖(版本根据你的Flink版本调整):
<dependency>
    <groupId>com.ververica</groupId>
    <artifactId>flink-connector-canal-kafka</artifactId>
    <version>3.2.0</version>
</dependency>
  1. 修改transactions表定义:
tEnv.executeSql(
        "create table transactions \n" +
                "(\n" +
                "account_id bigint,\n" +
                "amount bigint,\n" +
                "transaction_time timestamp(3),\n" +
                "watermark for transaction_time as transaction_time - interval '5' second\n" +
                ") with \n" +
                "(\n" +
                "'connector'='canal-kafka',\n" +
                "'topic'='transactions',\n" +
                "'properties.bootstrap.servers'='hadoop102:9092',\n" +
                "'scan.startup.mode'='latest-offset',\n" +
                "'canal-json.ignore-parse-errors'='true'\n" +
                ");"
);

临时应急方案:丢弃null记录(不推荐)

若需临时绕过问题,取消代码中注释的配置,启用丢弃逻辑:

import org.apache.flink.configuration.Configuration;

// ...

EnvironmentSettings env = EnvironmentSettings.inStreamingMode();
TableEnvironment tEnv = TableEnvironment.create(env);
Configuration config = tEnv.getConfig().getConfiguration();
config.setString("table.exec.sink.not-null-enforcer", "DROP");

此方案会丢弃无法解析的记录,仅作临时应急,不建议长期使用。

内容的提问来源于stack exchange,提问作者Eric_h9

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 23:15:53