第 14 章 · 集成生态 — 可视化演示

四个交互演示:① Kafka → MV → MergeTree 三件套 · ② MaterializedPostgreSQL 同步 · ③ 表函数对比卡 · ④ 实时数仓全链路

Kafka → MV → MergeTree 三件套架构

看一条业务消息从 Kafka 到 ClickHouse 落盘的完整路径。注意 events_kafka 表是虚拟的,本身不存数据:

📱 业务客户端
埋点 SDK / 微服务
🟫 Kafka topic: user_events
(等待消息...)
① events_kafka
ENGINE = Kafka(...)
「虚拟表 / 消费者」
不存数据
② MV 物化视图
SELECT ... FROM events_kafka
「搬运工」
INSERT 触发器
③ events_target
ENGINE = MergeTree
「真正存数据」
查询打这张表
关键洞察:三件套缺一不可。Kafka 表是「水龙头」(不存数据),MV 是「水管」(搬运),MergeTree 是「水池」(真正存)。 可以挂多个 MV 把同一份 Kafka 数据分流到不同结构的目标表(明细一份、聚合一份、错误一份)。

MaterializedPostgreSQL 同步流程

基于 PG 逻辑复制 slot 把整库同步到 CK 做 OLAP 镜像。看一条 PG 上的 INSERT 如何到达 CK:

📝 业务应用(OLTP 写)
INSERT INTO orders VALUES (1001, 'A', 99.0)
🐘 PostgreSQL 主库
wal_level = logical
表:orders / users / products
📜 WAL 日志(逻辑解码 pgoutput)
{op: c, table: orders, after: {...}}
🔌 Replication Slot: ck_slot_xxx
CK 注册的「订阅入口」
未消费的 WAL 不会被删
🐝 CK MaterializedPostgreSQL 引擎
周期性拉取 → 各 PG 表落地为 ReplacingMergeTree
🪞 CK 镜像表
learn_pg_mirror.orders / users / ...
OLAP 查询跑这里
生产坑提醒: CK 进程挂了一定要去 PG 上 SELECT pg_drop_replication_slot('ck_slot_xxx'), 否则 PG 会保留所有未消费的 WAL,几小时就能把磁盘撑爆!

外部表函数 / 引擎对比卡

ClickHouse 提供了一整套「把 X 当 CK 表来查」的能力,每张卡片对应一种数据源:

mysql()
表函数 · MySQL
现场拉 MySQL 表数据。WHERE 条件会下推到 MySQL。
SELECT * FROM mysql(
  'host:3306',
  'shop', 'orders',
  'user', 'pwd'
) WHERE id = 1001;
实时下推每查必连
postgresql()
表函数 · PostgreSQL
同上,但是连 PG。多了个 schema 参数(PG 三层结构)。
SELECT * FROM postgresql(
  'pg-host:5432',
  'shop', 'orders',
  'user', 'pwd', 'public'
);
实时下推支持 schema
s3()
表函数 · 对象存储
直接读 / 写 S3 / OSS / MinIO 上的 Parquet / CSV / ORC。多文件并行读,Parquet row group 跳过。
SELECT count(*) FROM s3(
  'https://b/events/*.parquet',
  'Parquet',
  'AKIA...', 'SECRET...'
);
数据湖多文件并行归档/读冷
hdfs()
表函数 · HDFS
读写 HDFS 上的文件。和 s3() 用法几乎一样,连 Hadoop 生态。
SELECT count(*) FROM hdfs(
  'hdfs://nn:9000/logs/*.json',
  'JSONEachRow'
);
Hadoop多格式
remote() / remoteSecure()
表函数 · CK 跨集群
临时连另一个 CK 节点 / 集群。常用于数据搬家、跨集群对账。
INSERT INTO new_db.events
SELECT * FROM remote(
  'old-ck:9000',
  old_db.events,
  'user','pwd'
);
免配置数据搬家跨集群
url()
表函数 · HTTP
把任意 URL 返回的 CSV / JSON 当成 CK 表查。常用于读公开数据集。
SELECT * FROM url(
  'https://x/data.csv',
  'CSVWithNames'
);
HTTPad-hoc

Kafka → 明细 → 物化视图聚合 → 看板:实时数仓全链路

典型「实时埋点统计」链路:从一条业务事件落到看板的完整旅程:

📱 App SDK 上报
user_id=1001
type=click
🟫 Kafka
(user_events)
12 partitions
events_kafka
CK Kafka 引擎
消费组:dwh_v1
MV1:搬运 + 解析
SELECT * FROM events_kafka
events_local(明细层)
MergeTree · 7d TTL · 亿级行
MV2:分钟聚合
GROUP BY toStartOfMinute, event_type
uniqExactState / countState
events_minute_agg
AggregatingMergeTree
状态聚合(小数据量)
📊 Grafana 实时看板
每分钟 PV/UV/独立设备数
设计精髓:明细层保留原始数据(可重放),聚合层用 AggregatingMergeTree 存「中间状态」。 看板查询不是 GROUP BY 海量明细,而是 uniqExactMerge(state) 把已经压缩过的草图合并 —— 这就是 CK 能秒级返回亿级数据聚合的根本。