You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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                                          |    

我的操作步骤如下:

  1. 创建仅含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
  1. 创建将值以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字段即可:

  1. 确认原始流定义(已创建可跳过):
CREATE STREAM CONVERT_FROM_FIREWALL ( message string) WITH (KAFKA_TOPIC='from-syslog-firewall', VALUE_FORMAT='JSON');
  1. 使用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;
  1. 验证结果:
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 11:23:31