KSQL Stream中CDC_TIMESTAMP字段带(string)前缀问题及解决咨询
问题原因
- 原Kafka Topic中的
CDC_TIMESTAMP字段是嵌套JSON结构({"string":"2023-07-19T16:43:44.926000000000"}),但创建CORTEX_TYPE流时,你将该字段定义为VARCHAR类型。 - KSQL会把整个嵌套JSON对象直接序列化为字符串,导致流中该字段值是完整的嵌套结构字符串,因此Sink到下游后会带有
"string":前缀。
解决方法
修改流的定义,提取嵌套结构里的时间值即可,以下两种方式任选:
方式1:通过STRUCT类型匹配嵌套结构
先定义源流时匹配原始数据的嵌套结构,再提取目标字段:
CREATE STREAM CORTEX_TYPE ( ACCOUNT_TYPE_KY VARCHAR KEY, ADDED_BY VARCHAR, CDC_TIMESTAMP STRUCT<string VARCHAR> ) WITH ( KAFKA_TOPIC='arb.snpdh_CTX_ACCOUNT_TYPE', VALUE_FORMAT='JSON' ); CREATE STREAM CORTEX_TYPE_STREAM WITH(VALUE_FORMAT='AVRO') AS SELECT ACCOUNT_TYPE_KY, ADDED_BY, CDC_TIMESTAMP->string AS CDC_TIMESTAMP FROM CORTEX_TYPE;
方式2:使用JSON提取函数直接处理
如果不想修改源流定义,可直接用EXTRACT_JSON_FIELD函数从字符串化的JSON中提取时间值:
CREATE STREAM CORTEX_TYPE_STREAM WITH(VALUE_FORMAT='AVRO') AS SELECT ACCOUNT_TYPE_KY, ADDED_BY, EXTRACT_JSON_FIELD(CDC_TIMESTAMP, '$.string') AS CDC_TIMESTAMP FROM CORTEX_TYPE;
内容的提问来源于stack exchange,提问作者Mohamed Ayman
相关产品推荐
相关产品推荐

