ClickHouse 实时数仓:一条消息从 Kafka 到查询结果的完整旅程
2026/8/31 13:52:05 网站建设 项目流程

ClickHouse 实时数仓:一条消息从 Kafka 到查询结果的完整旅程

【免费下载链接】ClickHouseClickHouse® is a real-time analytics database management system项目地址: https://gitcode.com/GitHub_Trending/cli/ClickHouse

ClickHouse 实时数仓是一种"一套库同时管住数据链路两端"的用法:实时写入和分析查询都在同一个系统里完成,不用为流和批各维护一条系统。下文跟着一条事件消息的旅程走一遍:它如何被毫秒级接入、列式落盘、在后台被预先聚合,最后被秒级查询取回;顺带给出三步可复现的搭建步骤,以及不该用 ClickHouse 的场景。

半夜那次等了 10 秒的查询

常见的故事是这样的:晚十一点,用户的下单事件陆续进入 Kafka,业务方要一份"每用户每小时下单数"的早间报表。行式数据库上跑一遍全量聚合,查询要十秒以上,数据量翻倍就超时;另建一套 Flink 流式聚合,则要自己维护状态存储、一致性保证和回刷逻辑。

这类需求的本质是:数据持续在产生,而你希望它一到就能查、查到就出数。

让"边写边查"成立的三个机制

先给结论:ClickHouse 的做法是"先写、后并、边聚合边查"。

列式存储 + 向量化执行。同一列连续落盘,查询只取 3 列时不会触碰另外 20 列;执行引擎按数据块(向量化,见 src/Processors/)批量处理,比逐行计算快得多。实时写入这边也不亏:写入先进内存再落成磁盘文件,写入延迟在毫秒级。

写入与合并分两阶段。每次写入生成一个独立数据分区(Part),查询直接扫 Part;后台另有线程池异步合并、压缩 Part,数据放得越久查询越快。这是 MergeTree 系存储引擎能扛高频小写入又保持查询快的原因,核心实现在 src/Storages/MergeTree/。

多种接入通道。流数据有 Kafka、NATS(JetStream)等消息队列表引擎;批数据有 S3 表引擎和 Iceberg 这类湖仓格式,实现见 src/Storages/ObjectStorage/DataLakes/Iceberg/。流和批都从同一个库进来。

从零搭链路的 3 步

第 1 步:建一张流表当"入口"。用 Kafka 表引擎声明消费哪个 topic、什么格式,ClickHouse 后台自动维护消费者,持续拉取消息:

CREATE TABLE kafka_events ( event_time DateTime, user_id UInt64, event_type String ) ENGINE = Kafka() SETTINGS kafka_broker_list = 'kafka:9092', kafka_topic_list = 'user_events', kafka_format = 'JSONEachRow';

第 2 步:用物化视图"顺路聚合"。物化视图不是视图,而是真表:每次向源表写入数据,聚合结果立刻写入目标表,实时数据天然"预聚合"好了,原理与实现见 src/Storages/MaterializedView/。

CREATE MATERIALIZED VIEW user_stats ENGINE = SummingMergeTree() ORDER BY (user_id, toDate(event_time)) AS SELECT user_id, toDate(event_time) AS d, count() AS cnt FROM kafka_events GROUP BY user_id, d;

第 3 步:直接查聚合表。早间报表跑一个小 SELECT 即可,不用碰原始消息;历史批数据在对象存储上(Iceberg/Parquet)时,直接读入并与实时表 JOIN。配置上建议调的两个参数,默认值见 programs/server/config.xml:

参数建议值作用
max_insert_threads8 或 auto大 INSERT 的并行度
background_pool_size16后台合并线程数,写入频繁时调大

坑与边界:这三种情况别硬用

最常见的坑是把 ClickHouse 当 OLTP 数据库:它不适合高频改单行的小事务。MergeTree 的 UPDATE/DELETE 不是即时操作,如果一张表的主要用途是"改",请用 ReplacingMergeTree 思路设计,或直接换库。

其次是规模与集群边界:单机性能很好,但数据量到 TB 级、QPS 很高时,需要 Distributed 引擎加 Keeper 做多分片,这套部署的门槛不低,建议先压测。

第三是预聚合粒度:物化视图只能回答建视图时定义的粒度内的查询,粒度太粗答不了细问题,太细表会膨胀。常见做法是分钟级和天级各建一个。

下一步可以深挖的方向

  • 冷热分离:用 src/Disks/ 的多磁盘策略把实时热数据放本地盘,用 TTL 把历史数据移到 S3,查询热、存储便宜。
  • 新版 Kafka2 引擎用 Keeper 存消费位点、支持分区与分片亲和,多分片部署值得看,变更记录见 CHANGELOG.md。
  • 想知道查询能跑多快,看 tests/performance/ 的基准定义与 docs/ 官方文档。

一句话:如果你的问题是"一条持续产生的数据流,想查得快、聚合好、一套系统搞定",这套流批一体的架构能帮你省掉另一套系统的建设和维护成本。

【免费下载链接】ClickHouseClickHouse® is a real-time analytics database management system项目地址: https://gitcode.com/GitHub_Trending/cli/ClickHouse

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询