嵌套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
相关产品推荐
相关产品推荐

