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

📅 发布时间:2026/8/31 13:52:09
ClickHouse 实时数仓:一条消息从 Kafka 到查询结果的完整旅程 ClickHouse 实时数仓一条消息从 Kafka 到查询结果的完整旅程【免费下载链接】ClickHouseClickHouse® is a real-time analytics database management system项目地址: https://gitcode.com/GitHub_Trending/cli/ClickHouseClickHouse 实时数仓是一种一套库同时管住数据链路两端的用法实时写入和分析查询都在同一个系统里完成不用为流和批各维护一条系统。下文跟着一条事件消息的旅程走一遍它如何被毫秒级接入、列式落盘、在后台被预先聚合最后被秒级查询取回顺带给出三步可复现的搭建步骤以及不该用 ClickHouse 的场景。半夜那次等了 10 秒的查询常见的故事是这样的晚十一点用户的下单事件陆续进入 Kafka业务方要一份每用户每小时下单数的早间报表。行式数据库上跑一遍全量聚合查询要十秒以上数据量翻倍就超时另建一套 Flink 流式聚合则要自己维护状态存储、一致性保证和回刷逻辑。这类需求的本质是数据持续在产生而你希望它一到就能查、查到就出数。让边写边查成立的三个机制先给结论ClickHouse 的做法是先写、后并、边聚合边查。列式存储 向量化执行。同一列连续落盘查询只取 3 列时不会触碰另外 20 列执行引擎按数据块向量化见 src/Processors/批量处理比逐行计算快得多。实时写入这边也不亏写入先进内存再落成磁盘文件写入延迟在毫秒级。写入与合并分两阶段。每次写入生成一个独立数据分区Part查询直接扫 Part后台另有线程池异步合并、压缩 Part数据放得越久查询越快。这是 MergeTree 系存储引擎能扛高频小写入又保持查询快的原因核心实现在 src/Storages/MergeTree/。多种接入通道。流数据有 Kafka、NATSJetStream等消息队列表引擎批数据有 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),仅供参考