这是《ClickHouse 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 实时链路的目标不是把所有数据尽快写入,而是在延迟、成本、准确性和可恢复性之间找到稳定平衡。
23.1 常见架构
直接消费:
业务 / 埋点
-> Kafka
-> ClickHouse Kafka 引擎
-> MergeTree
适合格式稳定、清洗简单、延迟要求不极端的场景。
流计算清洗:
Kafka raw
-> Flink
-> Kafka clean
-> ClickHouse
适合乱序、状态计算、维表补齐、多流 Join 和复杂异常处理。
混合链路:
明细 -> ClickHouse
汇总 -> Flink / ClickHouse 物化视图
检索 -> Elasticsearch
离线 -> 数据湖
不同系统承担不同语义,不强行用一个系统覆盖所有需求。
23.2 延迟定义
必须拆开指标:
| 指标 | 含义 |
|---|---|
| upstream delay | 业务发生到 Kafka |
| consume delay | Kafka 到处理系统 |
| write delay | 处理系统到 ClickHouse |
| merge delay | 合并或去重完成 |
| query delay | 查询可见 |
ReplacingMergeTree 数据通常很快可见,但物理去重可能滞后。报表如果要求最新状态,应查询 argMax 语义,而不是等待合并完成。
23.3 写入缓冲
推荐写入端微批:
最大行数:50000
最大间隔:2s
失败重试:指数退避
最大重试:有限次数
死信:本地文件或 Kafka retry topic
Flink 伪代码:
sink.bufferMaxRows = 50000;
sink.bufferInterval = Duration.ofSeconds(2);
sink.failurePolicy = RetryThenDeadLetter;
缓冲越大吞吐越高,但故障时的损失窗口也越大。金融、计费等关键数据要配合外部可重放机制。
23.4 乱序与水印
事件时间可能乱序:
10:00:03
10:00:01
10:00:02
处理原则:
- Flink 使用 watermark 和窗口;
- ClickHouse 使用事件日期分区,不用处理时间替代业务日期;
- 允许延迟数据写入历史分区;
- 报表支持迟到数据重跑;
- 状态流用版本表达顺序。
23.5 维表补齐
| 位置 | 特点 |
|---|---|
| 上游写入前 | 数据完整,但维表变更复杂 |
| Flink Lookup | 实时性好,需缓存和 TTL |
| ClickHouse 字典 | 查询时补齐,适合小维表 |
| 每日宽表 | 简单稳定 |
高频维度建议写入时固化;缓慢变化维度可用字典或版本化维表。
23.6 数据质量
数据质量关卡:
schema 校验
必填校验
类型转换
枚举值校验
时间合法性
主键唯一性
金额范围
去重规则
分区边界
治理结果表:
CREATE TABLE analytics.data_quality
(
check_time DateTime,
topic String,
rule_name String,
error_count UInt64,
sample String
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(check_time)
ORDER BY (check_time, topic, rule_name);
23.7 幂等设计
典型规则:
事件流:event_id + 去重窗口
状态流:business_key + version
指标流:窗口 + 维度 + 可重放幂等
常用策略:
- Kafka offset 记录;
- 批次 ID;
- 业务版本;
- 事件唯一 ID;
- 分区重写;
- 死信重放。
23.8 监控
关键链路指标:
| 层 | 指标 |
|---|---|
| 业务 | 事件量、成功率 |
| Kafka | lag、分区均衡、错误 |
| 计算 | 反压、checkpoint、重启 |
| ClickHouse | 写入行数、part、合并 |
| 查询 | 延迟、错误、扫描量 |
| 业务结果 | 核心指标对账 |
告警示例:
lag > 5 分钟
写入失败率 > 1%
目标表行数环比异常
too many parts 接近阈值
消费者消失
核心报表空结果
23.9 回压与降级
当 ClickHouse 压力过大:
- 写入端降低并发;
- 增大批次间隔;
- 暂停低优先级 topic;
- 依靠 Kafka 保留数据等待恢复;
- 查询限流;
- 关闭高成本报表;
- 保留死信和重放入口。
不要在未定位原因时直接清空队列或重启集群。先确认是 IO、CPU、内存、part 还是依赖故障。
本章小结
实时链路要把数据契约、缓冲写入、乱序处理、幂等、监控和降级一起设计。简单链路可以直接用 Kafka 引擎,复杂语义交给 Flink 等流计算层。ClickHouse 的职责是稳定批量落表和高效分析。
思考题
- 直接消费和流计算清洗如何选择?
- 事件延迟包含哪几段?
- 为什么要用事件时间分区?
- 维表补齐有哪些方案?
- 写入压力过大时如何降级?