ksqlDB解析带引号分隔消息遇SerializationException问题排查
问题
我有一个存储防火墙JSON格式消息的Kafka主题,每条消息包含一个值为逗号分隔字符串的message字段:
{ "@timestamp": 1690745699, "pri": "14", "time": "Jul 30 19:34:59", "ident": "FIREWALL", "message": "aa,bb,cc,dd" }
我希望将该字段拆分到另一个流,得到如下结果:
ksql> select * from CONVERT_FROM_PALOALTO_SER; +--------------------------------------------+--------------------------------------------+--------------------------------------------+--------------------------------------------+ |FIRST |SECOND |THIRD |FOURTH | +--------------------------------------------+--------------------------------------------+--------------------------------------------+--------------------------------------------+ |aa |bb |cc |dd |
我的操作步骤如下:
- 创建仅含
message字段的第一个流:
CREATE STREAM CONVERT_FROM_FIREWALL ( message string) WITH (KAFKA_TOPIC='from-syslog-firewall', VALUE_FORMAT='JSON');
查询结果正常:
ksql> select * from CONVERT_FROM_FIREWALL limit 1; +MESSAGE |+-|aa,bb,cc,dd
- 创建将值以
DELIMITED格式写入新主题的流:
CREATE STREAM CONVERT_FROM_FIREWALL_DELIMITED WITH (VALUE_FORMAT='DELIMITED', KAFKA_TOPIC='from-syslog-firewall_delimited') AS SELECT * FROM CONVERT_FROM_FIREWALL;
此时在Kafka中看到消息带有双引号。
3. 创建第三个流用于解析分隔符消息:
CREATE STREAM CONVERT_FROM_FIREWALL_SER ( first string, second string, third string, fourth string) WITH (KAFKA_TOPIC='from-syslog-firewall_delimited', VALUE_FORMAT='DELIMITED', VALUE_DELIMITER=',');
但新消息进入Kafka时,ksqlDB抛出错误:
Caused by: org.apache.kafka.common.errors.SerializationException: Column count mismatch on deserialization. topic: from-syslog-firewall_delimited, expected: 4, got: 1
补充:若向Kafka发送不带引号的消息,ksqlDB无报错且能正常解析。请问我哪里操作有误?该如何解决此问题?
问题原因与解决方案
错误原因
你创建CONVERT_FROM_FIREWALL_DELIMITED流时,直接SELECT *输出整个message字段,ksqlDB会把这个字符串字段序列化为带双引号的单值分隔符消息(比如"aa,bb,cc,dd")。当后续流尝试按逗号拆分出4个字段时,会把整个带引号的字符串当成1个字段,自然出现列数不匹配的错误。
最优解决方案:一步拆分字段
不需要中间的DELIMITED格式主题,直接在ksqlDB里拆分message字段即可:
- 确认原始流定义(已创建可跳过):
CREATE STREAM CONVERT_FROM_FIREWALL ( message string) WITH (KAFKA_TOPIC='from-syslog-firewall', VALUE_FORMAT='JSON');
- 使用
SPLIT函数拆分字段,直接创建目标流:
CREATE STREAM CONVERT_FROM_FIREWALL_SER AS SELECT SPLIT(message, ',')[1] AS FIRST, SPLIT(message, ',')[2] AS SECOND, SPLIT(message, ',')[3] AS THIRD, SPLIT(message, ',')[4] AS FOURTH FROM CONVERT_FROM_FIREWALL;
也可以用STRING_SPLIT_TO_ARRAY函数实现相同效果:
CREATE STREAM CONVERT_FROM_FIREWALL_SER AS SELECT ARRAY_ELEMENT(STRING_SPLIT_TO_ARRAY(message, ','), 1) AS FIRST, ARRAY_ELEMENT(STRING_SPLIT_TO_ARRAY(message, ','), 2) AS SECOND, ARRAY_ELEMENT(STRING_SPLIT_TO_ARRAY(message, ','), 3) AS THIRD, ARRAY_ELEMENT(STRING_SPLIT_TO_ARRAY(message, ','), 4) AS FOURTH FROM CONVERT_FROM_FIREWALL;
- 验证结果:
SELECT * FROM CONVERT_FROM_FIREWALL_SER LIMIT 1;
即可得到拆分后的字段结果。
备选方案:保留中间主题的处理方式
如果必须通过中间DELIMITED主题中转,需要确保写入中间主题的是拆分后的多个字段,而非原始带逗号的字符串:
-- 拆分后写入DELIMITED格式的中间主题 CREATE STREAM CONVERT_FROM_FIREWALL_DELIMITED WITH (VALUE_FORMAT='DELIMITED', KAFKA_TOPIC='from-syslog-firewall_delimited') AS SELECT SPLIT(message, ',')[1] AS FIRST, SPLIT(message, ',')[2] AS SECOND, SPLIT(message, ',')[3] AS THIRD, SPLIT(message, ',')[4] AS FOURTH FROM CONVERT_FROM_FIREWALL; -- 创建解析流(此时主题内是4个无引号的分隔值) CREATE STREAM CONVERT_FROM_FIREWALL_SER ( first string, second string, third string, fourth string) WITH (KAFKA_TOPIC='from-syslog-firewall_delimited', VALUE_FORMAT='DELIMITED', VALUE_DELIMITER=',');
内容的提问来源于stack exchange,提问作者Dmitry
相关产品推荐
相关产品推荐

