ksqlDB多列分组报错咨询:非键字段能否用于分组?
核心结论
可以使用流中任意字段(包括不属于事件键的字段)进行多列分组,并非根本无法实现,问题出在聚合后生成的表的键格式配置上。
报错原因分析
你最初遇到的错误:
Key format does not support schema. format: KAFKA schema:
Persistence{columns=[DATESTRING KEY,ACTIONSTRING KEY],
features=[]} reason: The 'KAFKA' format only supports a single field.
Got: [DATESTRING KEY,ACTIONSTRING KEY]
本质是:KEY_FORMAT='KAFKA'(原生Kafka字节键格式)仅支持单个键字段,但你按date, action多列分组后,生成的聚合表需要复合键来存储这两个分组字段。原流使用了KEY_FORMAT='KAFKA',ksqlDB默认会沿用这个格式创建聚合表,而它无法承载复合键,因此报错。
无需修改原主题结构的解决方案
不需要把分组字段提前放到事件键里,有两种直接的解决方式:
方案1:创建聚合表时显式指定支持复合键的键格式
在CREATE TABLE语句中通过WITH子句指定KEY_FORMAT为JSON、AVRO等支持复合键的格式即可:
CREATE TABLE TABLE_1 AS SELECT date, action, SUM(VALUE) as number FROM stream_1 GROUP BY date, action EMIT CHANGES WITH (KEY_FORMAT='JSON'); -- 或AVRO、PROTOBUF等支持复合键的格式
方案2:先对原流重分区,再聚合
如果后续需要基于复合键做更多操作,可以先将原流按分组字段重分区,生成一个带复合键的新流,再进行聚合:
-- 第一步:生成重分区流,将date和action设为复合键,指定键格式为JSON CREATE STREAM STREAM_1_REPARTITIONED AS SELECT date, action, userId, value FROM stream_1 PARTITION BY date, action EMIT CHANGES WITH (KEY_FORMAT='JSON'); -- 第二步:基于重分区流执行聚合 CREATE TABLE TABLE_1 AS SELECT date, action, SUM(VALUE) as number FROM STREAM_1_REPARTITIONED GROUP BY date, action EMIT CHANGES;
这种方式还能优化聚合性能,因为数据已经按分组字段分布到对应的分区中。
你修改后成功的原因
你将date和action移到事件键中,并设置KEY_FORMAT='JSON',本质是让原流的键本身就是复合键,聚合时沿用了支持复合键的格式,因此能正常执行。但这只是解决问题的一种方式,而非必须步骤。
内容的提问来源于stack exchange,提问作者slymore

