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

Flink消费Debezium变更日志JSON反序列化失败,请求排查DDL问题

问题

我正尝试使用Flink消费Debezium生成的变更事件日志,对应的JSON数据如下:

{
    "schema":{

    },
    "payload":{
        "before":null,
        "after":{
            "team_config_id":3800,
            "team_config_team_id":"team22bcb26e-499a-41e6-8746-b7d980e79e04",
            "team_config_sfdc_account_id":null,
            "team_config_sfdc_account_url":null,
            "team_config_business_type":5,
            "team_config_dpsa_status":0,
            "team_config_desc":null,
            "team_config_company_id":null,
            "team_config_hm_count_stages":null,
            "team_config_assign_credits_times":null,
            "team_config_real_renew_date":null,
            "team_config_action_date":null,
            "team_config_last_action_date":null,
            "team_config_business_tier_notification":"{}",
            "team_config_create_date":1670724933000,
            "team_config_update_date":1670724933000,
            "team_config_rediscovery_tier":0,
            "team_config_rediscovery_tier_notification":"{}",
            "team_config_sfdc_industry":null,
            "team_config_sfdc_market_segment":null,
            "team_config_unterminated_note_id":0
        },
        "source":{

        },
        "op":"c",
        "ts_ms":1670724933149,
        "transaction":null
    }
}

我尝试了两种方式声明输入Schema:

第一种方式是直接解析JSON数据:

create table team_config_source (
      `payload` ROW <
        `after` ROW <
          ...
          team_config_create_date timestamp(3),
          team_config_update_date timestamp(3),
          ...
        >
      >
    ) WITH (
    'connector' = 'kafka',
    ...
    'format' = 'json'
    )

但Flink抛出错误:org.apache.flink.formats.json.JsonToRowDataConverters$JsonParseException: Fail to deserialize at field: team_config_create_date,原因是java.time.format.DateTimeParseException: Text '1670724933000' could not be parsed at index 0。Flink不支持这种格式的时间戳吗?

我还尝试了另一种方式,使用内置的debezium格式:

create table team_config_source (
      team_config_create_id int,
      ...
    ) WITH (
    'connector' = 'kafka',
    ...
    'format' = 'debezium-json'
    )

但Flink又出现另一个错误:java.io.IOException: Corrupt Debezium JSON message,原因是java.lang.NullPointerException。我看到有人说更新事件的before值不能为null,但这条消息是创建事件。

请问有人能帮忙检查我的DDL吗?

解决方案

第一种JSON解析方式的问题处理

Flink的JSON格式默认仅支持解析字符串格式的时间戳(如"2022-12-10T10:15:33.000"),无法直接将数字类型的毫秒时间戳解析为timestamp类型,这就是报错的核心原因。可以通过以下两种方式解决:

  1. 先以数值类型读取,再转换为时间戳
    把时间字段定义为BIGINT,查询时通过TO_TIMESTAMP_LTZ函数转换为时间类型:

    create table team_config_source (
      `payload` ROW <
        `after` ROW <
          team_config_id INT,
          team_config_team_id STRING,
          -- 其他字段省略...
          team_config_create_date BIGINT,
          team_config_update_date BIGINT
        >
      >
    ) WITH (
      'connector' = 'kafka',
      'topic' = 'your-topic',
      'properties.bootstrap.servers' = 'your-brokers',
      'format' = 'json'
    );
    
    -- 查询转换示例
    SELECT 
      TO_TIMESTAMP_LTZ(payload.after.team_config_create_date, 3) AS create_date,
      TO_TIMESTAMP_LTZ(payload.after.team_config_update_date, 3) AS update_date
    FROM team_config_source;
    
  2. 自定义JSON格式的时间解析规则(Flink 1.16+支持)
    在WITH参数中添加配置,开启数字转时间戳的功能:

    create table team_config_source (
      `payload` ROW <
        `after` ROW <
          -- 字段定义同前
          team_config_create_date TIMESTAMP(3),
          team_config_update_date TIMESTAMP(3)
        >
      >
    ) WITH (
      'connector' = 'kafka',
      ...
      'format' = 'json',
      'json.timestamp-format.standard' = 'ISO-8601',
      'json.number-as-timestamp' = 'true' -- 开启数字时间戳解析
    );
    

    注意:此配置仅在Flink 1.16及以上版本可用,且需确保时间戳为毫秒级(与你的数据匹配)。

第二种Debezium-JSON格式的问题处理

报错NullPointerException并非因为before为null(创建事件的before本身就是null,属于合法情况),而是DDL字段与Debezium消息的after字段不匹配(字段名错误、类型不匹配),或未正确配置Debezium格式参数。

正确的DDL写法如下:

create table team_config_source (
  team_config_id INT,
  team_config_team_id STRING,
  team_config_sfdc_account_id STRING,
  team_config_sfdc_account_url STRING,
  team_config_business_type INT,
  team_config_dpsa_status INT,
  team_config_desc STRING,
  team_config_company_id STRING,
  team_config_hm_count_stages INT,
  team_config_assign_credits_times INT,
  team_config_real_renew_date BIGINT,
  team_config_action_date BIGINT,
  team_config_last_action_date BIGINT,
  team_config_business_tier_notification STRING,
  team_config_create_date BIGINT,
  team_config_update_date BIGINT,
  team_config_rediscovery_tier INT,
  team_config_rediscovery_tier_notification STRING,
  team_config_sfdc_industry STRING,
  team_config_sfdc_market_segment STRING,
  team_config_unterminated_note_id INT
) WITH (
  'connector' = 'kafka',
  'topic' = 'your-topic-name',
  'properties.bootstrap.servers' = 'your-kafka-brokers',
  'format' = 'debezium-json',
  'debezium-json.ignore-parse-errors' = 'false',
  'debezium-json.schema-include' = 'false' -- 你的消息中schema为空,关闭schema解析
);

关键注意事项:

  • 字段名必须与after中的字段完全一致(大小写敏感),比如你之前写的team_config_create_id是错误的,应为team_config_id
  • 时间字段仍需先定义为BIGINT,再通过TO_TIMESTAMP_LTZ转换为时间类型
  • 添加debezium-json.schema-include = 'false',避免解析空的schema字段引发异常
  • 确保字段类型与消息中的数据类型匹配,比如字符串用STRING,整数用INT或BIGINT

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 20:25:17