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

ClickHouse Kafka表JSON解析报错(错误码26):Kafka表与物化视图创建方案咨询

解决ClickHouse Kafka表JSON解析错误(错误码26)的方案

咱们先拆解下你遇到的问题:错误码26是ClickHouse的JSON解析失败异常,核心原因是你用的20.9.3版本比较旧,它对Nested类型与JSON数组的映射解析逻辑不够健壮,导致无法正确识别data字段对应的数组结构,才抛出了“预期左引号”的错误。

下面给你两个靠谱的解决方案,你可以根据需求选择:

方案1:先接收原始JSON字符串,再在物化视图中解析(最稳妥)

这种方式绕开了旧版本ClickHouse对Nested类型直接解析的bug,先把Kafka消息原封不动存成字符串,再通过物化视图解析成目标结构,容错性更强。

步骤1:创建Kafka原始数据接收表

CREATE TABLE queue_raw (
    raw_data String
) ENGINE = Kafka 
SETTINGS 
    kafka_broker_list = 'xxx,xxx,xxx',
    kafka_topic_list = 'librenms_flow',
    kafka_group_name = 'niop-billing-test-group',
    kafka_format = 'RawBLOB', -- 直接接收原始字符串,避免提前解析出错
    kafka_num_consumers = 1;

步骤2:创建最终存储表

用来保存解析后的数据,用MergeTree引擎保证查询性能:

CREATE TABLE queue_final (
    deviceId Int8,
    data Nested(
        port_id String,
        time Int64,
        in Float32,
        out Float32
    )
) ENGINE = MergeTree()
ORDER BY deviceId;

步骤3:创建物化视图实现自动解析

物化视图会自动消费queue_raw里的原始数据,解析后写入queue_final:

CREATE MATERIALIZED VIEW queue_mv TO queue_final AS
SELECT
    JSONExtractInt(raw_data, 'deviceId') AS deviceId,
    -- 先把data字段的数组提取成原始JSON字符串数组
    JSONExtractArrayRaw(raw_data, 'data') AS data_raw,
    -- 遍历数组,逐个解析每个元素的字段,映射成Nested类型
    arrayMap(
        x -> tuple(
            JSONExtractString(x, 'port_id'),
            JSONExtractInt(x, 'time'),
            JSONExtractFloat(x, 'in'),
            JSONExtractFloat(x, 'out')
        ),
        data_raw
    ) AS data
FROM queue_raw;

方案2:调整Kafka表的格式参数(快速尝试)

如果你不想加中间表,可以试试给原Kafka表添加几个兼容旧版本的参数,修复解析逻辑:

CREATE TABLE queue (
    deviceId Int8,
    data Nested(
        port_id String,
        time Int64,
        in Float32,
        out Float32
    )
) ENGINE = Kafka 
SETTINGS 
    kafka_broker_list = 'xxx,xxx,xxx',
    kafka_topic_list = 'librenms_flow',
    kafka_group_name = 'niop-billing-test-group',
    kafka_format = 'JSONEachRow',
    kafka_num_consumers = 1,
    kafka_json_allow_single_quotes = 0, -- 强制只解析双引号的JSON,避免格式混淆
    kafka_skip_broken_messages = 1; -- 临时跳过解析失败的消息,避免阻塞消费

不过这个方案依赖旧版本的兼容性,万一Kafka里有格式异常的消息,还是容易出问题,所以更推荐方案1。

验证方法

创建完成后,你可以执行下面的语句验证数据是否正常写入:

SELECT * FROM queue_final LIMIT 10; -- 方案1用这个
-- 方案2直接查Kafka表:SELECT * FROM queue LIMIT 10;

如果能查到数据,说明解析正常。要是还有问题,可以查看Kafka消费状态排查:

SELECT * FROM system.kafka_consumers WHERE table = 'queue_raw'; -- 对应方案1的表名

内容的提问来源于stack exchange,提问作者Yanpeng

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 12:38:15