如何用Flink SQL将JSON_VALUE读取的字符串转为TIMESTAMP_LTZ并解决字段为null
问题解决:Upsert-Kafka表TIMESTAMP_LTZ字段转换失败导致值为Null
问题背景
已创建Upsert-Kafka表roles_created_raw_v1,表定义如下:
CREATE TABLE roles_created_raw_v1 ( id VARCHAR, created VARCHAR PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'sink_topic', 'properties.bootstrap.servers' = 'localhost:29092,localhost:39092', 'properties.group.id' = 'sink_topic_id', 'value.format' = 'json', 'key.format' = 'json', 'properties.allow.auto.create.topics' = 'true', 'value.json.timestamp-format.standard' = 'ISO-8601', 'sink.parallelism' = '3' );
执行插入语句时,created字段在Kafka主题和表中均为Null:
insert into roles_created_raw_v1 select JSON_VALUE(contentJson, '$.id') as id, to_timestamp(JSON_VALUE(contentJson, '$.created'), 'yyyy-MM-ddTHH:mm:ss.SSSZ') as created from some_raw_table;
已知contentJson的内容为:
"contentJson": "{\"created\":\"2023-02-04T04:12:07.925Z\"}"
且使用JSON_VALUE(contentJson, '$.created' RETURNING STRING)可以正常获取到字符串格式的时间值,需解决将该字符串正确转换为TIMESTAMP_LTZ类型的问题。
解决方案
1. 修改表结构,匹配目标字段类型
原表中created定义为VARCHAR,与目标TIMESTAMP_LTZ类型不符,需调整表结构:
CREATE TABLE roles_created_raw_v1 ( id VARCHAR, created TIMESTAMP_LTZ(3) -- 对应时间字符串的毫秒精度 PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'upsert-kafka', 'topic' = 'sink_topic', 'properties.bootstrap.servers' = 'localhost:29092,localhost:39092', 'properties.group.id' = 'sink_topic_id', 'value.format' = 'json', 'key.format' = 'json', 'properties.allow.auto.create.topics' = 'true', 'value.json.timestamp-format.standard' = 'ISO-8601', 'sink.parallelism' = '3' );
2. 使用适配时区的转换函数
原to_timestamp函数适用于无时区的时间字符串,对于带Z(UTC时区)的ISO-8601格式,应使用to_timestamp_ltz函数,且无需指定格式模板(Flink会自动识别标准ISO-8601格式):
insert into roles_created_raw_v1 select JSON_VALUE(contentJson, '$.id' RETURNING STRING) as id, to_timestamp_ltz(JSON_VALUE(contentJson, '$.created' RETURNING STRING), 3) as created from some_raw_table;
补充说明
TIMESTAMP_LTZ(3)中的3需与输入时间字符串的毫秒精度保持一致;- 显式指定
JSON_VALUE的返回类型RETURNING STRING,可避免隐式类型转换带来的异常; - 表配置中
value.json.timestamp-format.standard已设为ISO-8601,Flink会自动将TIMESTAMP_LTZ类型序列化为标准JSON时间字符串写入Kafka。
内容的提问来源于stack exchange,提问作者hitesh
相关产品推荐
相关产品推荐

