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

嵌套JSON数组结构的Kafka Stream创建与KsqlDB字段提取问题

解决方案

步骤1:重新定义流结构匹配JSON嵌套格式

不要将case1定义为字符串类型,直接映射为和JSON结构匹配的嵌套STRUCT+ARRAY类型,省去手动解析JSON的成本:

CREATE STREAM Stream1 (
    case1 STRUCT<
        case2 ARRAY<
            STRUCT<
                case3A STRUCT<var1 VARCHAR, var2 VARCHAR, var3 VARCHAR>,
                case3B STRUCT<var4 VARCHAR, var5 VARCHAR, var6 VARCHAR>,
                case3C STRUCT<var7 VARCHAR, var8 VARCHAR, var9 VARCHAR>
            >
        >
    >
) WITH (KAFKA_TOPIC='testtopic', VALUE_FORMAT='JSON', PARTITIONS=6);

如果实际字段类型不是字符串,可根据真实业务场景调整STRUCT内的字段类型定义。

步骤2:用EXPLODE函数展开数组获取全量元素

KsqlDB内置的EXPLODE函数就是你需要的类explode能力,会把数组的每个元素拆分为独立的行,配合LATERAL TABLE即可实现需求:

SELECT 
    item->case3A->var1,
    item->case3B->var4,
    item->case3C->var7
FROM Stream1
LATERAL TABLE EXPLODE(case1->case2) AS t(item);

如果使用的是较旧版本的ksqlDB,可省略LATERAL TABLE,直接写为:

SELECT 
    item->case3A->var1,
    item->case3B->var4,
    item->case3C->var7
FROM Stream1, EXPLODE(case1->case2) AS t(item);

兼容原有流定义的方案

如果你不想修改已创建的Stream1定义,保持case1为VARCHAR类型也可以实现,先解析出数组再展开即可:

SELECT 
    EXTRACTJSONFIELD(item, '$.case3A.var1') as var1,
    EXTRACTJSONFIELD(item, '$.case3B.var4') as var4
FROM (
    SELECT EXPLODE(EXTRACTJSONFIELD(case1, '$.case2')) as item
    FROM Stream1
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 11:06:05