ksqlDB创建Stream获取Kafka消息Header报错,求解决方案
KSQL创建流时HEADERS语法错误的解决方法
错误根源
你写的CREATE STREAM语句里,在字段列表中加HEADERS关键字的写法不符合KSQL的语法规范,这是导致报错的直接原因——KSQL并不支持在字段定义里这样声明headers字段。
正确实现方式
要在KSQL中获取Kafka消息的key、value和headers,有两种常用的正确写法:
方法1:创建流时开启headers支持,查询时获取所有headers
先创建包含key和value的基础流,开启headers支持:
CREATE STREAM A ( k1 VARCHAR KEY, v1 VARCHAR ) WITH ( KAFKA_TOPIC = 'ratings', VALUE_FORMAT = 'delimited', HEADERS = 'INCLUDE' );
之后查询时,用HEADERS()函数就能拿到所有headers的数组:
SELECT k1, v1, HEADERS() AS h1 FROM A EMIT CHANGES;
方法2:直接在查询中提取指定header
如果只需要某一个特定的header,不用在创建流时额外配置,直接在查询里提取:
SELECT k1, v1, HEADER('h1') AS h1_value FROM A EMIT CHANGES;
注意事项
- 确认你的KSQL版本在5.3及以上,更早的版本不支持headers操作
- 检查JMeter发送的消息确实携带了headers,不然查询时会返回null
内容的提问来源于stack exchange,提问作者Ruban Srinivasan
相关产品推荐
相关产品推荐

