ClickHouse 实时流处理实战:从 Kafka 接入到集群部署与调优
发布时间:2026/9/28 6:34:52来源:尧图网络
实时数据这块有个老矛盾数据从 Kafka 里源源不断流进来计算层跑得飞快结果到了查询端要么是 Elasticsearch 聚合慢得让人抓狂要么是传统数仓 T1 根本接不住分钟级窗口。我这两年把 ClickHouse 放到实时流处理链路里专职做分析侧的存储与查询引擎算是把这个瓶颈彻底治好了。这篇文章就把我在这套架构上的完整实践——从引擎选型、链路设计、Kafka 接入、集群部署再到那些让人头大的重启报错和性能调优——一次性梳理清楚适合正在做实时数仓选型、或者已经被 ClickHouse 部署和流式接入折腾过几次的同学参考。1. 实时场景的引擎选型为什么ClickHouse能扛住流处理的下游压力1.1 一次真实的选型对比去年做一个用户行为实时分析平台上游是埋点日志经 Kafka 汇聚要求是事件发生后 30 秒内能看到聚合报表的秒级延迟。第一版我用了 MySQL 加定时任务结果数据量一到千万级聚合查询直接飙到十几秒业务方天天催。后来把 MySQL 换成了 Elasticsearch单条明细查询倒是快了但多维度 GROUP BY 和漏斗分析依然吃力而且内存吃得吓人。当时摆在面前的可选方案还有 Druid 和 ClickHouse。Druid 的架构对实时摄入支持很成熟但部署组件多、学习曲线陡。ClickHouse 则是一个纯列式 OLAP 引擎单机性能极强SQL 语法贴近习惯而且摄入数据就是标准的 INSERT和 Kafka、Flink 对接非常直接。我最后选了 ClickHouse原因就三句话列式存储天然适合宽表聚合分布式表屏蔽了集群复杂度物化视图能把实时预聚合的成本摊到写入路径上。1.2 ClickHouse 与传统实时数仓的定位差异很多人把 ClickHouse 当成快一点的 MySQL或者省内存的 Elasticsearch这个理解其实有偏差。ClickHouse 的核心竞争力来自三件事列式存储、向量化执行、稀疏索引。列式存储意味着查询只需要读取涉及的列而不是像行式存储那样把整行数据都捞出来。在行为分析这种字段多、聚合频繁的场景里同样是计算一个 UV列式引擎读取的数据量可能只有行式的十分之一。向量化执行则是把 SQL 操作转换成对 CPU SIMD 指令友好的批处理循环单核吞吐能跑到每秒数亿行。稀疏索引配合分区裁剪让它在十亿行级别的表上做带时间范围的聚合时依然可以稳定在百毫秒到秒级。和 Flink、Spark Streaming 这类流计算引擎不同ClickHouse 不做状态管理和事件时间语义它更擅长的是存和查。所以在实时链路里它和流引擎是互补关系Flink 负责复杂的窗口计算、状态维护和事件时间处理ClickHouse 负责把结果以及明细高效落盘并提供极速查询。这个定位分清楚之后整个架构的职责边界就清晰了。2. 端到端实时链路Kafka、Flink 与 ClickHouse 的角色分工2.1 整体架构与数据流向我常用的实时链路是经典的三段式数据源埋点、业务库 binlog、物联网设备上报→ Kafka流处理层Flink 或轻量 Kafka Consumer做清洗、关联、窗口聚合结果写入 ClickHouse由报表服务和即席查询直接读取Kafka 在这条链路里扮演的是削峰填谷的缓冲带。流量洪峰来了上游可以猛写Flink 消费不过来就攒着ClickHouse 不会被瞬时写入压垮。ClickHouse 单表写入性能虽然强但也不是无限吸收的有了 Kafka 这层缓冲消费速率就能被控制在合理区间。Flink 层解决的是数据形态转换原始埋点日志可能是嵌套 JSON需要展开成宽表不同来源的数据需要按用户 ID 关联事件需要按窗口求 PV/UV。这些逻辑如果全压在 ClickHouse 的 SQL 里既难维护又会拖慢查询所以我的经验是能在 Flink 里算完的聚合就不要留给查询时再算。2.2 链路中 ClickHouse 的双表设计我在实践中发现实时链路里只建一张明细表是不够的。推荐的做法是同时维护两张表明细表以原始事件粒度存储字段全、保留期短用于问题回溯和下钻聚合表由物化视图实时维护以分钟或小时为粒度预聚合用于报表查询这两张表共用同一个 Kafka 数据源但物化视图只把聚合结果写进聚合表。查询端优先走聚合表只有需要看单条明细时才查明细表。这个设计和 Lambda 架构的精简版很像但没有批处理和流处理两套代码的维护负担因为 ClickHouse 的物化视图是增量更新的天然把实时写入和实时聚合合并成了一条路径。2.3 为什么不在 ClickHouse 里直接做复杂关联需要说句实话ClickHouse 的 JOIN 能力正在变强但复杂多表关联仍然不是它的主场。我在项目里一般只允许两类 JOIN一是小表维表几万行以内用字典或 Global JOIN二是同分布键的大表 JOIN。跨大表的多级关联宁可放到 Flink 里用状态编程完成也别指望 ClickHouse 的 SQL 层替你扛。这也顺带解释了为什么实时链路里 Flink 不是可有可无的。真正的分工是Flink 做整形ClickHouse 做存储与加速Kafka 做缓冲。三者各管一段出了问题也好定位——数据不对先查 Flink 的日志查询慢再查 ClickHouse 的执行计划链路清晰太多。3. 数据接入实操Kafka 引擎表与物化视图的配合细节3.1 用 Kafka 引擎表完成最轻量的写入如果你的实时计算逻辑不复杂只有简单的过滤和字段重命名可以不用 Flink直接靠 ClickHouse 的 Kafka 引擎表完成接入。CREATE TABLE kafka_user_events ( user_id UInt64, event_type String, event_time DateTime, properties String ) ENGINE Kafka() SETTINGS kafka_broker_list 192.168.1.10:9092,192.168.1.11:9092, kafka_topic_list user_events, kafka_group_name clickhouse_consumer_group, kafka_format JSONEachRow, kafka_num_consumers 4;这段建表语句是接入表它本身不存数据只是消费 Kafka 消息的通道。真正落数据的是后面的目标表再通过物化视图把两者串起来。我刚开始用的时候踩过一个坑直接在 Kafka 引擎表上做 SELECT发现能查到数据以为万事大吉结果一重启服务数据全没了。Kafka 引擎表本质上是流式缓冲区不是存储必须通过物化视图把数据搬运到 MergeTree 系引擎的表里。3.2 物化视图把流数据变成预聚合结果下面是核心的创建顺序-- 1. 先建目标明细表 CREATE TABLE user_events ( user_id UInt64, event_type String, event_time DateTime, properties String ) ENGINE MergeTree() PARTITION BY toYYYYMMDD(event_time) ORDER BY (event_time, user_id); -- 2. 再建聚合目标表 CREATE TABLE user_events_minute_agg ( event_time DateTime, event_type String, uv UInt64, pv UInt64 ) ENGINE SummingMergeTree() PARTITION BY toYYYYMMDD(event_time) ORDER BY (event_time, event_type); -- 3. 创建物化视图把 Kafka 接入表的数据同时转发明细表和聚合表 CREATE MATERIALIZED VIEW mv_events_to_agg TO user_events_minute_agg AS SELECT toStartOfMinute(event_time) AS event_time, event_type, uniqExact(user_id) AS uv, count() AS pv FROM kafka_user_events GROUP BY toStartOfMinute(event_time), event_type;这里有几个细节值得展开说。第一物化视图的 SELECT 是对 Kafka 表做实时查询每来一批消息就增量算一次结果写入 TO 子句指定的目标表。它是增量式的所以不要担心全量重算的问题。第二uniqExact在数据量大时很费内存我在几十亿行级别的 PV/UV 场景里会换成uniqCombined它是有损但误差极小的近似算法性能能提升一个数量级。报表场景完全够用审计场景才需要精确值。第三聚合表用 SummingMergeTree当同一排序键分钟事件类型的多个批次到达时后台会自动把相同键的行合并求和。但注意SummingMergeTree 的合并是后台异步的查询时可能遇到同一个 key 有多行的情况。解决方式是查询时仍要带上GROUP BY event_time, event_type或者用sum配合聚合函数来兜底。3.3 写入语义与格式选型Kafka 消息格式我推荐 JSONEachRow 或 ClickHouse 原生的 Native 格式。JSON 可读性好但解析开销大压测下来吞吐大概只有 Native 的一半。如果链路里有 Flink 做转换直接从 Flink 侧以 Native 格式批量写入最划算。批量写入的另一个关键是攒批。我一般设置攒批策略积攒 5000 条或者 2 秒超时二者先到先发。这样既避免了一条条 INSERT 带来的高频提交也不会因为攒太多导致数据延迟超过 SLA。ClickHouse 对大批量 INSERT 的优化非常激进一次写 5 万行比 5 千行快得多。注意ClickHouse 的副本间同步是异步的。如果写入的是一张 ReplicatedMergeTree 表写入方只需要等一个副本确认即可返回。务必在业务侧对主备延迟有预期否则你做主备切换测试时会发现数据短暂不一致这是设计如此不是故障。4. 生产环境部署集群配置、系统参数与重启故障排查4.1 集群拓扑与副本方案先给出一套我在生产环境验证过的部署参数3 节点规模配置项推荐值说明CPU16 核以上ClickHouse 是 CPU 密集型核数比内存重要内存64GB 起建议至少配置 1:4 的 CPU 与内存比磁盘NVMe SSD x2RAID1日志和元数据放系统盘数据目录独立挂载操作系统Ubuntu 22.04 / Rocky Linux 9内核版本不低于 5.x 即可副本机制ReplicatedMergeTree clickhouse-keeper用内置 keeper 替代 ZooKeeper省心很多人在 Rocky Linux 9 上装 ClickHouse 会遇到依赖问题主要是缺libstdc较新版本或unixODBC。官方 RPM 仓库里其实已经解决了依赖只要先执行yum install -y unixODBC再装clickhouse-server和clickhouse-client就能避免大部分报错。4.2 重启报错 failed to flush system log already exists 的根因这个报错很多人在网上搜过我先说结论它通常不是数据损坏而是非正常停机导致 system 数据库中的日志表出现了孤儿元数据。现象是重启后 server 起不来日志里反复出现类似failed to flush system log already exists的信息。我第一次遇到时以为是磁盘坏了排查了半天才发现是/var/lib/clickhouse/metadata/system目录下的表元数据和实际表状态不一致。处理办法分两步先确认数据目录里是否残留旧的system相关目录。ClickHouse 的 system 库query_log、query_thread_log、part_log等是内存表加异步落盘的如果上次进程被 kill -9临时文件可能没有清干净。把/var/lib/clickhouse/metadata/system/下的残留.sql文件移走备份再启动服务。等服务稳定后执行SYSTEM FLUSH LOGS重建日志表。如果这种方法仍无法解决还有一个更保守的路径将config.xml里的system_log段落注释掉先让服务起来再逐步恢复日志表。但这只是应急手段长期运行还是要有正常停机的流程。4.3 正常停机和异常崩溃的处理习惯我在生产环境里踩过几次坑后形成了一套固定流程日常维护停服先SYSTEM STOP MERGES再SYSTEM FLUSH LOGS然后用clickhouse-server stop正常退出。保证日志表和合并任务都有序收尾。异常崩溃后的恢复先检查system.errors表如果还能起来的话再看/var/log/clickhouse-server/clickhouse-server.err.log末尾的堆栈信息。切忌直接删目录重启那会把整个表结构弄丢。升级版本前务必备份metadata目录和config.d下的自定义配置。ClickHouse 小版本升级大体兼容跨大版本比如 22.x 到 23.x一定要先读 changelog 里的兼容性说明。提示任何线上改动配置后执行clickhouse-server --config/etc/clickhouse-server/config.xml --validate-config先做语法校验。这个命令能抓到一半以上的配置错误别等启动失败才后悔。4.4 集群部署策略从单机到多副本的扩展路径我的建议是先单机、后集群。不要一开始就搭复杂的多副本集群先把单机跑通确认查询性能和写入吞吐符合预期再按需加节点。单机转集群时两步走把表引擎换成ReplicatedMergeTree配置 clickhouse-keeper 提供元数据协调。这里要重点设计ORDER BY键副本表同一分区的数据在节点间是按主键范围分布的选错排序键会导致某些节点数据量严重倾斜。创建Distributed分布式表作为查询入口应用层只写分布式表配合internal_replication true让副本之间的复制由表引擎自己负责应用无需感知节点健康状态。副本数并不是越多越好。我记得有一次为了高可用上了 3 副本结果节点间同步占用大量带宽高峰期的写入延迟反而上升了。实时分析场景对读多写少更友好两个副本完全够用一个承担实时写入一个承担报表查询互相还能分摊压力。5. 实时查询调优分区裁剪、主键设计与常见慢查询治理5.1 分区键和时间字段的配合在实时流处理场景里绝大多数查询都带时间范围。因此分区键选对收益几乎是立竿见影的。事件流表我建议按天分区PARTITION BY toYYYYMMDD(event_time)。如果数据量日均过亿再缩短到小时级分区。分区粒度太细也有副作用——Part 数量过多会拖慢后台合并严重时触发too many parts报错。经验法则是单分区数据量至少保持几十万行以上分目录数量控制在几千以内。主键ORDER BY字段则要结合最常见的查询条件。行为分析里通常是某用户在某时间段做了哪些事所以我会把(event_time, user_id)作为排序键。注意 ClickHouse 的排序键不要求唯一它只决定数据在磁盘上的物理顺序。对用户维度的下钻查询把user_id放在event_time前面反而更快因为相同用户的数据会聚在相邻块里稀疏索引跳过的块更多。这个选择要按实际查询频率权衡。5.2 用 EXPLAIN 和小表数据演习定位慢查询遇到慢查询不要瞎调参先看执行计划EXPLAIN PLAN SELECT event_type, count() FROM user_events WHERE event_time today() GROUP BY event_type;重点看两处ReadType是不是FULL全表扫描Parts数量是否为 0没命中分区裁剪。如果WHERE条件里的时间字段没有走分区多半是数据类型不一致——比如表里存的是DateTime查询却传了字符串隐式转换会让索引失效。我还有个习惯把慢 SQL 缩小到一个分区内的数据做演习人为制造一个只含百万行的临时表在临时表上反复调整PREWHERE和GROUP BY顺序。这样能立刻看出瓶颈是在 IO 扫描还是 CPU 聚合不至于在生产上试错。5.3 内存参数与并发控制实时报表的特点是查询频率高、单次扫描数据量可控。这种情况下把max_threads适当调低反而能提升并发吞吐。我跑过对比并发 30 个查询时max_threads8比默认的 16 表现更好因为线程过多会触发大量上下文切换。最关键的参数是max_memory_usage。默认配置下单个查询能用完整机内存这在生产环境是灾难。我一般设置为物理内存的 60%再配合max_bytes_before_external_group_by打开内存溢出的外部聚合路径这样即使 GROUP BY 数据量超出内存也不会 OOM只是会慢一点。6. 踩坑实录从数据积压到磁盘爆满的复盘6.1 物化视图消费迟滞导致的 Kafka 积压我曾经遇到过一次很诡异的数据迟到明细表数据正常增长但聚合表的指标停留在 20 分钟前。排查过程复盘给你参考。第一步用SELECT * FROM system.mutations看有没有堆积的 mutation——没有。第二步查system.processes发现一个来自物化视图的 INSERT 查询长时间处于Waiting for query execution状态。第三步看到后台的 Merge 任务正在跑一个大分区的合并把 IO 占满了物化视图的写入排队排在后面。这就是 ClickHouse 的一个特性物化视图是同步阻塞式消费 Kafka 的。如果目标表所在的 MergeTree 在进行大合并写入会等待。解决方案是限制合并任务的并发度merge_tree max_part_merging_threads2/max_part_merging_threads replicated_max_merges_in_queue4/replicated_max_merges_in_queue /merge_tree同时把聚合表的休眠时间old_parts_lifetime从默认的 8 分钟降到 1 分钟避免过期分区占着磁盘空间。这套组合拳下来积压问题基本消失。6.2 磁盘写满前的静默故障还有一个坑比上面的更隐蔽ClickHouse 在磁盘写满时不会立刻崩溃而是先拒绝新的 INSERT表现是 Flink 侧不停报写入超时但 ClickHouse 日志里只有零星的Disk is full警告。如果你只看 Flink 的监控很容易以为是网络问题。我现在给所有集群都加了磁盘告警阈值设在 80%。同时把storage_policy配成多路径卷系统盘和 SSD 数据盘分开放system日志与用户数据。另外别忘了TTL策略明细表保留 7 天聚合表保留 90 天到期自动删除。没有 TTL 的实时表迟早会因为磁盘问题在凌晨三点叫你起来。6.3 一个建议把 ClickHouse 版本锁在稳定大版本内ClickHouse 迭代太快社区版每个月一个大版本。我的经验是生产环境锁定在某个大版本的最后一个 patch比如 23.8 LTS 系列新功能可以等两个版本稳定后再评估升级。流处理链路和 OLAP 引擎耦合紧密千万别为了一个新语法去追最新版稳定性永远是第一位。在 Flink 连接器版本选择上我会选官方维护的 Flink ClickHouse Connector 或自研的 HTTP 写入器。社区里很多第三方连接器停更严重对 ClickHouse 新版本不兼容最常见的问题就是写入的数据类型映射错误。如果有条件写一个基于 HTTP 接口的批量写入器并不难我后来的项目就换成了自研写入器单次批次 5 万行、50 毫秒超时跑了半年没出过问题。7. 最后聊几句实战体会这套 ClickHouse 加实时流处理的架构跑了快两年最大的体会是选型本身不难难的是把每个环节的边界划分清楚。Kafka 管缓冲Flink 管计算ClickHouse 管存储查询谁也别越界出问题的时候能够迅速定位到具体环节。具体到操作层面我建议新手从Kafka 引擎表 物化视图 单机 MergeTree起步先把链路跑通再逐步引入 Flink 和副本集群。千万别一开始就追求组件大而全否则你会在 ZooKeeper、keeper、分布式表、副本同步这些概念里绕晕反而忘了核心目标是让数据尽快可查。最后分享一个小技巧给每张实时表都建一个数据新鲜度监控查询用SELECT max(event_time) FROM table定期比对当前时间。这个简单的监控比任何复杂的指标面板都早发现积压问题。当最大值长时间停在几分钟前你基本可以断定链路某个环节卡住了直接按上面说的排查路径去看就行。
网站建设高端定制企业官网