这是《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 客户端版本及设置相关。
必须监控:
- consumer group lag;
- topic 分区数;
- 消费者数量;
- 消息格式错误;
- ClickHouse 写入异常;
- 目标表行数。
查看消费者:
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 管理,且版本兼容性要提前验证。生产建议:
- 上游发送标准化字段;
- 时间统一时区;
- 金额用字符串或定点格式;
- 必填字段稳定;
- 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 保持批量落表和高效分析。
思考题
- Kafka 引擎表保存长期数据吗?
- 为什么需要物化视图?
- 消费者数量如何确定?
- 如何处理重复消息?
- lag 增长时按什么顺序排查?