KSQL基于Kafka Topic的JSON键值创建流报错问题咨询
问题根因排查
1. ARRAY类型报错原因
KSQL中定义数组类型时必须明确指定数组内存储的元素类型,你语句中仅写了array未声明元素类型,因此触发未知类型报错。从你的示例数据看cw字段是整数数组,需要定义为ARRAY<INT>。
2. 查询timestamp返回null原因
你对KSQL流字段的映射逻辑存在误解:
- 你在流定义中额外新增的
ts struct<payload struct<timestamp bigint>>字段没有任何对应数据源,KSQL默认会将流字段映射到Kafka消息value的顶级字段,你的消息value中根本没有叫ts的顶级字段,因此查询返回null。 - 你在测试只定义key字段时能查到value的属性,是因为KSQL对于未定义的value字段会默认做隐式解析,不需要显式声明ts字段,timestamp本身就在value的payload结构下,直接通过
value->payload->timestamp调用即可。
解决方案
方案1:全量显式定义字段(推荐,性能更稳定)
注意带特殊符号、和关键字重名的字段需要用反引号包裹:
CREATE STREAM s_devices ( -- 声明该字段映射到Kafka消息的Key,而非Value `key` STRUCT<payload STRUCT<sourceName STRING, jobName STRING>> KEY, -- value的完整结构定义,数组明确指定元素类型为INT `value` STRUCT< payload STRUCT< fields STRUCT<cw ARRAY<INT>>, `timestamp` BIGINT, expires STRING, `connection-name` STRING > > ) WITH ( KAFKA_TOPIC='devices', VALUE_FORMAT='JSON', KEY_FORMAT='JSON', -- 可选配置:如果要将消息中的timestamp作为流的事件时间,可添加该行 TIMESTAMP='value->payload->timestamp' );
方案2:简化定义(依赖隐式解析,适合快速调试)
如果不需要全量字段,也可以只定义需要的Key字段,value字段走默认隐式解析即可:
CREATE STREAM s_devices ( `key` STRUCT<payload STRUCT<sourceName STRING>> KEY ) WITH ( KAFKA_TOPIC='devices', VALUE_FORMAT='JSON', KEY_FORMAT='JSON' );
验证查询
SELECT `key`->payload->sourceName AS source_name, `value`->payload->`timestamp` AS device_ts, `value`->payload->fields->cw AS cw_array, `value`->payload->`connection-name` AS connection_name FROM s_devices EMIT CHANGES LIMIT 10;
内容的提问来源于stack exchange,提问作者Al3
相关产品推荐
相关产品推荐

