ClickHouseNotes

第 22 章:Kafka 引擎

zjc 于 2026-01-22 发布

这是《ClickHouse 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 表引擎让 ClickHouse 直接消费 Kafka 消息,通常配合物化视图写入 MergeTree 表。它是实时链路入口,不是长期存储。

22.1 架构

Kafka topic
    |
Kafka engine table
    |
Materialized View
    |
MergeTree / ReplacingMergeTree

Kafka 表负责读取消息,物化视图负责转换,目标表负责存储和查询。

22.2 创建表

CREATE TABLE analytics.kafka_events
(
    event_time DateTime,
    event_id String,
    user_id UInt64,
    event_type LowCardinality(String),
    city_id UInt32,
    amount Decimal64(2)
)
ENGINE = Kafka
SETTINGS
    kafka_broker_list = 'kafka-01:9092,kafka-02:9092',
    kafka_topic_list = 'events',
    kafka_group_name = 'clickhouse-analytics',
    kafka_format = 'JSONEachRow',
    kafka_row_delimiter = '\n',
    kafka_num_consumers = 4,
    kafka_thread_per_consumer = 1,
    kafka_max_block_size = 65536,
    kafka_skip_broken_messages = 10;

常用设置:

设置 说明
kafka_broker_list Kafka 地址
kafka_topic_list topic,可多个
kafka_group_name 消费组
kafka_format 消息格式
kafka_num_consumers 消费者数量
kafka_max_block_size 批大小
kafka_skip_broken_messages 跳过坏消息数量

22.3 目标表和视图

CREATE TABLE analytics.events_local
(
    event_date Date,
    event_time DateTime,
    event_id String,
    user_id UInt64,
    event_type LowCardinality(String),
    city_id UInt32,
    amount Decimal64(2)
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, city_id, user_id);
CREATE MATERIALIZED VIEW analytics.events_mv
TO analytics.events_local
AS
SELECT
    toDate(event_time) AS event_date,
    event_time,
    event_id,
    user_id,
    event_type,
    city_id,
    amount
FROM analytics.kafka_events;

22.4 消费位点

ClickHouse 使用 Kafka consumer group 保存位点。位点提交语义、重平衡和可见性与 Kafka 客户端版本及设置相关。

必须监控:

  1. consumer group lag;
  2. topic 分区数;
  3. 消费者数量;
  4. 消息格式错误;
  5. ClickHouse 写入异常;
  6. 目标表行数。

查看消费者:

SELECT database, table, consumer_id, assignments.topic
FROM system.kafka_consumers;

22.5 消费者数量

有效消费者数量不能超过 topic 分区数:

topic partitions = 12
kafka_num_consumers = 4

增加 consumers 不一定提升吞吐。还要确认 ClickHouse 资源、目标表写入批次、Kafka 拉取配置、网络带宽和单分区乱序影响。

22.6 数据格式

JSONEachRow 示例:

{"event_time":"2026-08-25 10:00:00","event_id":"E001","user_id":1001,"event_type":"pay","city_id":1,"amount":99.00}

Avro、Protobuf 需要相关格式配置和 schema 管理,且版本兼容性要提前验证。生产建议:

  1. 上游发送标准化字段;
  2. 时间统一时区;
  3. 金额用字符串或定点格式;
  4. 必填字段稳定;
  5. schema 演进向后兼容。

22.7 异常消息

kafka_skip_broken_messages 只跳过有限坏消息,不是通用容错方案。

更稳妥的链路:

Kafka raw topic
   |
解析服务 / Flink
   |
valid topic + dead letter topic
   |
ClickHouse

如果直接消费 raw topic,建议至少保存原始消息、错误原因、topic、partition、offset、处理时间和重试状态。

22.8 管理操作

停止和恢复消费:

SYSTEM STOP CONSUMERS analytics.kafka_events;
SYSTEM START CONSUMERS analytics.kafka_events;

修改 Kafka 表通常需要删除重建。操作顺序:

1. 停止写入或记录位点
2. 停止 consumers
3. 删除物化视图
4. 删除 Kafka 表
5. 重建 Kafka 表
6. 重建物化视图
7. 校验无重复无缺失

22.9 幂等与重复

Kafka 至少一次投递常见,重放会造成重复。

目标表 方案
MergeTree 事件 ID 查询去重或接受重复
ReplacingMergeTree 业务键 + 版本
CollapsingMergeTree 正负状态记录
SummingMergeTree 指标可累加且重放可控

对账:

SELECT event_id, count()
FROM analytics.events_local
WHERE event_date = today()
GROUP BY event_id
HAVING count() > 1;

22.10 排查

问题 排查
无数据 group 位点、topic、格式、视图
lag 增长 分区数、消费者、写入压力
坏消息多 schema、字段类型、分隔符
重复 重放或视图重建
缺失 位点越界、坏消息跳过
消费停止 system.kafka_consumers、日志

本章小结

Kafka 引擎把消费、转换和存储串成一条链路。生产关键是格式契约、消费位点、坏消息处理、幂等写入和 lag 监控。复杂清洗和乱序处理放在流计算层,ClickHouse 保持批量落表和高效分析。

思考题

  1. Kafka 引擎表保存长期数据吗?
  2. 为什么需要物化视图?
  3. 消费者数量如何确定?
  4. 如何处理重复消息?
  5. lag 增长时按什么顺序排查?