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类型,这就是报错的核心原因。可以通过以下两种方式解决:
先以数值类型读取,再转换为时间戳
把时间字段定义为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;自定义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

