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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 19:57:44