如何创建以STRUCT为ARRAY元素类型的KSQL Table
实现包含STRUCT数组元素的KSQL Table思路
需求说明
将输入流中相同PRODUCT_SKU的商品变体数据聚合,生成以PRODUCT_SKU为键、包含变体STRUCT数组的KSQL Table。
输入流示例
{ "GTIN": "4066136261834", "PRODUCT_SKU": "UC_55.08_W08_900_black", "SIZE_NAME": "34", "CALC_SALES_PRICE": 199.90, "SALES_PRICE": 199.90 }
{ "GTIN": "4066136261846", "PRODUCT_SKU": "UC_55.08_W08_900_black", "SIZE_NAME": "36", "CALC_SALES_PRICE": 199.90, "SALES_PRICE": 199.90 }
期望输出KTable示例
{ "PRODUCT_SKU":"UC_55.08_W08_900_black", "variants":[ { "GTIN":"4066136261834", "SIZE_NAME":"34", "CALC_SALES_PRICE":199.90, "SALES_PRICE":199.90 }, { "GTIN":"4066136261846", "SIZE_NAME":"36", "CALC_SALES_PRICE":199.90, "SALES_PRICE":199.90 } ] }
实现步骤
1. 创建输入流
先为输入数据定义KSQL流,明确各字段数据类型:
CREATE STREAM product_variants_stream ( GTIN STRING, PRODUCT_SKU STRING, SIZE_NAME STRING, CALC_SALES_PRICE DOUBLE, SALES_PRICE DOUBLE ) WITH ( KAFKA_TOPIC='your-input-topic-name', -- 替换为实际Kafka主题名 VALUE_FORMAT='JSON', PARTITIONS=6 -- 根据实际场景调整分区数 );
2. 聚合生成目标KTable
通过GROUP BY按PRODUCT_SKU分组,用STRUCT()构造变体结构体,再用COLLECT_LIST()将结构体收集为数组:
CREATE TABLE product_variants_table AS SELECT PRODUCT_SKU, COLLECT_LIST( STRUCT( GTIN := GTIN, SIZE_NAME := SIZE_NAME, CALC_SALES_PRICE := CALC_SALES_PRICE, SALES_PRICE := SALES_PRICE ) ) AS variants FROM product_variants_stream GROUP BY PRODUCT_SKU;
关键说明
STRUCT()函数用于将多个字段组合成结构化对象,KSQL支持在聚合逻辑中直接构造STRUCT。COLLECT_LIST()会将分组内的所有STRUCT元素收集为数组,保留元素的原始插入顺序。- 若输入流的消息键不是
PRODUCT_SKU,建议在创建流时通过KEY_FORMAT和KEY参数指定键为PRODUCT_SKU,避免聚合时的数据重分区,提升处理性能。
内容的提问来源于stack exchange,提问作者RealTimeDataNerd
相关产品推荐
相关产品推荐

