主题
第 14 章 集成生态:让 ClickHouse 和上下游打通
学习目标:能用
Kafka引擎 + 物化视图搭建一条「Kafka → ClickHouse」实时摄入链路;理解MaterializedPostgreSQL引擎做 PG 全库 OLAP 镜像的原理;掌握mysql()/postgresql()/s3()/hdfs()/remote()这一套表函数的用法,把 ClickHouse 当成跨库联邦查询入口;知道 Airbyte / Debezium 这些 ETL 工具在数据架构里的位置。
0. 开篇:ClickHouse 不是孤岛
┌────────────────────────────────────────────────────────────────────┐
│ │
│ 业务库(OLTP) 消息总线 对象存储 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ MySQL │ ─binlog─→ │ Kafka │ ←──────→│ S3 / OSS │ │
│ │ Postgre │ Debezium │ RocketMQ │ CSV/ │ HDFS │ │
│ │ MongoDB │ │ Pulsar │ Parquet │ │ │
│ └────┬─────┘ └─────┬────┘ └─────┬────┘ │
│ │ │ │ │
│ │ MaterializedPG │ Kafka 引擎 │ s3() 表函数 │
│ │ mysql() 表函数 │ + MV 三件套 │ S3() 引擎 │
│ ▼ ▼ ▼ │
│ ┌──────────────────────────────────┐ │
│ │ ClickHouse 集群 │ │
│ │ 存全量明细 + 物化视图聚合 │ │
│ └────────────┬─────────────────────┘ │
│ │ │
│ │ remote() / Distributed │
│ ▼ │
│ ┌──────────────────────────────────┐ │
│ │ BI / 看板 / API 服务 │ │
│ │ Superset / Metabase / Grafana │ │
│ └──────────────────────────────────┘ │
└────────────────────────────────────────────────────────────────────┘ClickHouse 的产品哲学是:「我专心做查询引擎,数据怎么进来你随便」。所以官方提供了一大堆 引擎 与 表函数 让你把上下游接进来:
| 上下游 | 接入方式 | 典型场景 |
|---|---|---|
| Kafka / RocketMQ / Pulsar | Kafka 引擎 + MV | 实时摄入埋点、日志、CDC |
| MySQL | MySQL 引擎 / mysql() 表函数 | 跨库 join、维表 |
| PostgreSQL | PostgreSQL 引擎 / postgresql() 表函数 | 同上 |
| PostgreSQL(实时同步) | MaterializedPostgreSQL 引擎 | 整库做 OLAP 镜像 |
| S3 / OSS / MinIO | S3 引擎 / s3() 表函数 | 数据湖直读 Parquet |
| HDFS | HDFS 引擎 / hdfs() 表函数 | Hadoop 生态联邦 |
| 任意 CK 集群 | remote() / clusterAllReplicas() | 跨集群联邦查询 |
| 通用 ETL | Airbyte / Debezium / Flink CDC | 工具化批量入仓 |
📌 首次术语解释 · 表函数(Table Function):CK 提供的「SELECT 之中可以当成表来用的函数」。例如
SELECT * FROM s3('https://x/y/*.parquet', 'Parquet')会现场建一个临时虚拟表指向 S3 上的文件,无需先CREATE TABLE。
1. Kafka 引擎:实时摄入主力
1.1 三件套架构
ClickHouse 接 Kafka 的标准模式不是「Kafka 引擎一个表搞定」,而是 「Kafka 表 + MergeTree 表 + 物化视图」三件套:
┌────────────────────────────────────────────────────────────────────┐
│ │
│ ① Kafka 引擎表(消费者) │
│ ENGINE = Kafka(...) │
│ ─ 它是一个虚拟表,每次 SELECT 会从 Kafka 拉一批消息然后【消失】 │
│ ─ 不能 ORDER BY、不能存数据,单纯是「会喷数据的水龙头」 │
│ │ │
│ │ ② MV(物化视图)作为搬运工 │
│ │ MV = INSERT 触发器 │
│ ▼ │
│ ③ MergeTree 目标表(真正存数据) │
│ ENGINE = MergeTree │
│ ─ 应用查询数据时只查这张表 │
│ ─ Kafka 表的存在感为零 │
└────────────────────────────────────────────────────────────────────┘📌 为什么要分三张表? 因为 Kafka 引擎表本身不存数据(每条消息只能消费一次),必须把消费下来的数据「灌」到一张真正的表里。物化视图就是这条灌水管线 —— 每当 Kafka 引擎表「冒出」新数据,MV 自动把它 INSERT 到目标表。
1.2 完整 DDL 三件套
sql
-- ========== ① Kafka 引擎表(消费者)==========
CREATE TABLE learn_ck.events_kafka
(
event_time DateTime,
user_id UInt64,
event_type String,
properties String
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka1:9092,kafka2:9092,kafka3:9092',
kafka_topic_list = 'user_events',
kafka_group_name = 'ck_consumer_v1', -- 消费组(同组不会重复消费)
kafka_format = 'JSONEachRow', -- 消息格式
kafka_num_consumers = 4, -- 并发消费者数(≤ 该 topic 的分区数)
kafka_max_block_size = 65536, -- 一次取多少条
kafka_skip_broken_messages = 100; -- 容错:脏数据上限
-- ========== ③ MergeTree 目标表(真正存数据)==========
CREATE TABLE learn_ck.events_target
(
event_time DateTime,
user_id UInt64,
event_type LowCardinality(String),
properties String,
inserted_at DateTime DEFAULT now() -- 入库时间,方便排查延迟
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (event_type, user_id, event_time);
-- ========== ② 物化视图(搬运工)==========
CREATE MATERIALIZED VIEW learn_ck.mv_events_kafka_to_target
TO learn_ck.events_target -- ⭐ 注意:把数据写到 target 表
AS
SELECT
event_time,
user_id,
event_type,
properties
FROM learn_ck.events_kafka;注意 MV 的 TO 目标表 语法 —— 这是 CK 物化视图的关键差异,详见第 10 章。
1.3 关键 SETTINGS 详解
| 参数 | 含义 | 推荐值 |
|---|---|---|
kafka_broker_list | Kafka 集群地址 | 至少 3 个 broker |
kafka_topic_list | 订阅的 topic(可多个,逗号分隔) | 业务对应 |
kafka_group_name | Kafka 消费组名 | 每个独立消费链路一个 |
kafka_format | 消息格式 | JSONEachRow / CSV / Avro / Protobuf |
kafka_num_consumers | 同一节点的消费者线程数 | ≤ topic 的分区数 |
kafka_max_block_size | 一次拉取多少条凑成一个 Part | 65536(默认 8192 偏小) |
kafka_skip_broken_messages | 跳过多少条解析失败的消息 | 100(生产可改更大或 0 死循环) |
kafka_thread_per_consumer | 是否每个消费者一个线程 | 1(默认 0) |
kafka_commit_every_batch | 每批 commit offset | 0(默认按时间 commit) |
kafka_handle_error_mode | 错误处理:default / stream | stream 把错误也送到 MV |
1.4 偏移量、消费失败排查
1.4.1 看消费组当前位置
sql
SELECT
database, table, consumer_id,
assignments.topic, assignments.partition_id, assignments.current_offset,
last_poll_time, num_messages_read,
last_used,
rdkafka_stat, -- librdkafka 统计 JSON(信息量大)
last_exception
FROM system.kafka_consumers
FORMAT Vertical;1.4.2 常见故障
| 现象 | 原因 | 解法 |
|---|---|---|
| 消费 stuck 不动 | MV 写下游表失败(schema 错、磁盘满) | 看 system.errors 与 text_log;修 MV 后会自动追 |
| 消费速度极慢 | kafka_num_consumers 设得比分区数小 | 增大到 ≤ 分区数 |
大量 Cannot parse 报错 | 消息格式不一致 | 配 kafka_skip_broken_messages + 写一个错误链路 MV |
| Offset 漂移 | 消费组名变了 | 消费组名固定,慎改 |
| 重启后重复消费 | kafka_commit_every_batch=1 没开 + crash 在 commit 前 | 接受少量重复 + 在目标表用 ReplacingMergeTree 去重 |
1.4.3 暂停 / 拉起消费
sql
-- 暂停消费(不删表)
DETACH TABLE learn_ck.events_kafka;
-- 重新拉起
ATTACH TABLE learn_ck.events_kafka;
-- 重置 offset(要先 DETACH,再删 ZK 节点 / 改 group_name)1.5 错误链路:脏数据也能进 ClickHouse
把错误的消息也写入一张 errors 表,方便回溯:
sql
SET kafka_handle_error_mode = 'stream';
-- 在 Kafka 表上加两个虚拟列 _error / _raw_message
CREATE MATERIALIZED VIEW learn_ck.mv_events_errors
ENGINE = MergeTree ORDER BY (now(), _error)
AS
SELECT now() AS ts, _error AS err, _raw_message AS raw
FROM learn_ck.events_kafka
WHERE length(_error) > 0;这样脏消息既不阻塞主链路,又留下了痕迹。
2. MaterializedPostgreSQL:把 PG 整库做成 OLAP 镜像
2.1 一句话定义
MaterializedPostgreSQL引擎让 ClickHouse 通过 PG 的 逻辑复制 slot 实时订阅一个或多个 PG 表,把 PG 当主库 / CK 当只读分析镜像。
2.2 架构图
PostgreSQL(主) ClickHouse(CK)
┌──────────────┐ ┌────────────────────────┐
│ users / order│ → WAL → 逻辑复制slot │
│ products / ...│ (pgoutput) │ MaterializedPostgreSQL │
│ │ ← REPLICATION 协议→│ 引擎周期性拉取 → 落地 │
│ │ │ 各 PG 表为 ReplacingMT │
└──────────────┘ │ ── orders │
│ ── users │
│ ── products │
└────────────────────────┘
应用读 CK 跑 OLAP,
写仍走 PG(只读镜像)2.3 PG 端先开两件事
sql
-- 1) wal_level = logical
ALTER SYSTEM SET wal_level = 'logical';
-- 重启 PG
-- 2) 给同步用户开 REPLICATION 权限
CREATE USER ck_repl WITH REPLICATION LOGIN PASSWORD 'xxxxx';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO ck_repl;2.4 CK 端建库
sql
CREATE DATABASE learn_pg_mirror
ENGINE = MaterializedPostgreSQL(
'pg-host:5432', -- PG 地址
'source_db', -- PG 库名
'ck_repl', 'xxxxx' -- PG 同步账号
)
SETTINGS
materialized_postgresql_tables_list = 'orders,users,products', -- 订阅哪些表
materialized_postgresql_schema = 'public',
materialized_postgresql_max_block_size = 65536;CK 会在 PG 端自动创建一个 replication slot(名字形如 ck_<random>),并周期性拉取 WAL 解码后的变更。
2.5 注意事项
- 每张同步过来的表底层引擎是
ReplacingMergeTree,按 PG 主键去重。 - DDL 变更需要重建:CK 不会自动跟随 PG 的
ALTER TABLE。 - WAL slot 不释放会撑爆 PG 磁盘:CK 挂了一定要去 PG 上
pg_drop_replication_slot()清理。 - 不支持 PG 的所有类型:JSONB、数组、范围类型可能映射成 String,自定义类型不支持。
📌 生产用法:MaterializedPostgreSQL 适合「业务库 → OLAP 镜像」,不适合做 ETL 的中间层。复杂转换还是走 Debezium → Kafka → CK。
3. MySQL / PostgreSQL 表函数:跨库 ad-hoc 查询
3.1 mysql() 表函数
把任意 MySQL 表「现场虚拟成 CK 的表」用,每次查询都是去 MySQL 实时拉。
sql
SELECT user_id, count() AS pv
FROM mysql(
'mysql-host:3306', -- MySQL 地址
'shop', -- MySQL 库名
'orders', -- MySQL 表名
'reader_user', 'pwd' -- MySQL 账号 / 密码
)
WHERE created_at >= '2026-04-01'
GROUP BY user_id
ORDER BY pv DESC
LIMIT 10;性能小技巧:CK 会把过滤条件下推到 MySQL 执行,所以 WHERE created_at >= ... 不会把全表拉到 CK 再过滤。
sql
-- 也可以注册成持久化的 MySQL 引擎表,重复使用更方便
CREATE TABLE shop_orders_mysql
(
id BIGINT, user_id BIGINT, amount Decimal(18,2), created_at DateTime
)
ENGINE = MySQL('mysql-host:3306', 'shop', 'orders', 'reader_user', 'pwd');3.2 postgresql() 表函数
sql
SELECT *
FROM postgresql(
'pg-host:5432',
'shop',
'orders',
'reader_user', 'pwd',
'public' -- schema(PG 三层结构)
)
WHERE status = 'paid'
LIMIT 100;3.3 跨库 JOIN 实战
CK 真正强大的地方在于:可以把 MySQL、PG、CK 自己的表混在一个 SQL 里 JOIN:
sql
-- 用 CK 的事实表(events,亿级)
-- JOIN MySQL 的维度表(users,百万级,做字典)
SELECT
e.event_type,
u.country,
count() AS pv,
uniqExact(e.user_id) AS uv
FROM learn_ck.events_all e
LEFT JOIN mysql('mysql-host:3306', 'shop', 'users', 'r', 'p') u
ON e.user_id = u.id
WHERE e.event_time >= today() - 7
GROUP BY e.event_type, u.country
ORDER BY pv DESC
LIMIT 20;⚠ 性能提醒:MySQL 维表如果很大,别这样直接 JOIN,应该 预先全量同步到 CK 内做 Dictionary(详见第 9 章)。表函数适合「ad-hoc」「数据量适中(< 千万)」的场景。
4. S3 / HDFS:直接读写对象存储和数据湖
4.1 s3() 表函数:读 / 写 Parquet/CSV/ORC
sql
-- 读:跑一个 S3 上的 Parquet 文件夹的 SUM
SELECT count(), sum(amount)
FROM s3(
'https://my-bucket.s3.amazonaws.com/events/2026/04/*.parquet',
'Parquet',
'AKIA...', 'SECRET...' -- AccessKey / SecretKey(也可走 IAM Role)
);
-- 写:把 CK 的查询结果归档到 S3(Parquet 列存 + ZSTD 压缩)
INSERT INTO FUNCTION s3(
'https://my-bucket.s3.amazonaws.com/archive/2026-04.parquet',
'Parquet'
)
SELECT * FROM learn_ck.events_local
WHERE event_time >= '2026-04-01' AND event_time < '2026-05-01';支持的格式:Parquet / ORC / CSV / CSVWithNames / TSV / JSON* / Native 等。
4.2 S3 引擎表:定义为持久化表
sql
CREATE TABLE archive_2026
(
event_time DateTime, user_id UInt64, ...
)
ENGINE = S3(
'https://my-bucket.s3.amazonaws.com/archive/*.parquet',
'Parquet'
);
-- 跑查询时走 S3,CK 内部自动并行读多个文件
SELECT count() FROM archive_2026 WHERE event_time >= '2026-04-15';4.3 HDFS 引擎
sql
CREATE TABLE hdfs_logs
(
ts DateTime, msg String
)
ENGINE = HDFS('hdfs://nn:9000/logs/*.json', 'JSONEachRow');
-- 也有 hdfs() 表函数
SELECT count() FROM hdfs('hdfs://nn:9000/logs/*.json', 'JSONEachRow');📌 数据湖典型用法:用
s3()把 Spark / Flink 输出的 Parquet 归档读出来与 CK 实时数据 JOIN,构建「热数据 + 冷数据」的混合查询层。
5. remote() 表函数:跨集群即席查询
remote() 让你不需要建 Distributed 表,就能临时访问另一个 CK 集群(甚至单机):
sql
-- 单点远程
SELECT count() FROM remote('other-ck-host:9000', learn_ck.events_local, 'user', 'pwd');
-- 集群(按 ck_cluster_v2 配置 fan-out)
SELECT count() FROM remote('ck_cluster_v2', learn_ck.events_local);
-- 也有 remoteSecure()(TLS)/ cluster() / clusterAllReplicas()
SELECT count() FROM clusterAllReplicas('ck_cluster', learn_ck.events_local);
-- 区别:cluster() 每分片随机选一个副本;clusterAllReplicas() 所有副本都查(看一致性)最常见的两个用途:
- 数据搬家:
INSERT INTO new_cluster.events_all SELECT * FROM remote('old_cluster', ...) - 数据回查 / 比对:跨集群对账,看新老集群结果是否一致。
6. ETL 工具栈:Airbyte / Debezium / Flink CDC
CK 自带的引擎能解决大部分场景,但当数据源类型多、字段映射复杂时,专门的 ETL 工具更省心。
6.1 Debezium:CDC 黑话的代名词
CDC(Change Data Capture):把数据库的「写操作流」捕获出来,变成可消费的事件流。
Debezium 是 RedHat 开源的 CDC 中间件,支持几乎所有主流 OLTP(MySQL binlog / PG WAL / MongoDB oplog / SQL Server)。典型链路:
MySQL ─binlog→ Debezium Connector ─→ Kafka ─→ ClickHouse Kafka 引擎 ─→ MergeTreeCK 端:用 Kafka 引擎消费 Debezium 输出的格式(JSONEachRow 或 Avro),用 ReplacingMergeTree 按业务主键去重 + 用 __deleted 列处理删除事件。
6.2 Airbyte:低代码 EL 框架
Airbyte 提供 300+ Source / Destination 连接器,UI 拖拽就能配 MySQL → ClickHouse / Salesforce → ClickHouse 这种链路。适合数据团队人手不足、又要接很多业务源的场景。
6.3 Flink CDC:流式数仓的现代选择
Flink CDC(自 2.x 起含集成的 ClickHouse Sink)能做:
- 多源 CDC 合流后再写 CK
- 流式 JOIN / 窗口聚合后写 CK 物化视图层
- 端到端 exactly-once 语义
| 工具 | 优势 | 劣势 |
|---|---|---|
| CK Kafka 引擎 | 无依赖、性能好、配置简单 | 仅支持已落 Kafka 的数据 |
| Debezium | 业界事实标准、CDC 完整 | 需独立部署、需要 Kafka |
| Airbyte | 连接器多、UI 友好 | 性能不极致、部分 Source 是批量 |
| Flink CDC | 端到端流处理、复杂转换强 | 重 / 学习曲线陡 |
| Spark 写 CK | 适合大批量历史回灌 | 流式实时性弱 |
7. 真实案例:Kafka → MV → AggregatingMergeTree 的实时数仓链路
需求:实时统计「每个事件类型 × 每分钟」的 PV、UV、独立设备数。
7.1 整体架构
业务客户端 ─→ Kafka topic: user_events (JSON)
│
▼
┌────────────── ClickHouse ──────────────┐
│ │
│ events_kafka (Kafka 引擎,消费者) │
│ │ │
│ │ MV1:搬运 + 解析 │
│ ▼ │
│ events_local (MergeTree,明细,TTL 7d) │
│ │ │
│ │ MV2:聚合到分钟粒度 │
│ ▼ │
│ events_minute_agg │
│ (AggregatingMergeTree,状态聚合) │
│ │
└──────────────┬───────────────────────────┘
│ SELECT
▼
Grafana 看板(每秒刷新)7.2 完整 SQL
sql
-- 1. 明细层 events_local(详见第 5、7 章)
CREATE TABLE learn_ck.events_local (
event_time DateTime,
user_id UInt64,
device_id UInt64,
event_type LowCardinality(String),
properties String
)
ENGINE = MergeTree
PARTITION BY toYYYYMMDD(event_time)
ORDER BY (event_type, event_time, user_id)
TTL event_time + INTERVAL 7 DAY;
-- 2. Kafka 消费者
CREATE TABLE learn_ck.events_kafka (
event_time DateTime,
user_id UInt64,
device_id UInt64,
event_type String,
properties String
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka1:9092,kafka2:9092',
kafka_topic_list = 'user_events',
kafka_group_name = 'ck_realtime_dwh',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4;
-- 3. MV1:Kafka 表 → 明细表
CREATE MATERIALIZED VIEW learn_ck.mv_kafka_to_events
TO learn_ck.events_local AS
SELECT * FROM learn_ck.events_kafka;
-- 4. 聚合层 events_minute_agg(AggregatingMergeTree)
CREATE TABLE learn_ck.events_minute_agg (
minute DateTime,
event_type LowCardinality(String),
pv AggregateFunction(count),
uv AggregateFunction(uniqExact, UInt64),
devices AggregateFunction(uniqExact, UInt64)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMMDD(minute)
ORDER BY (event_type, minute);
-- 5. MV2:明细 → 分钟聚合
CREATE MATERIALIZED VIEW learn_ck.mv_events_to_minute
TO learn_ck.events_minute_agg AS
SELECT
toStartOfMinute(event_time) AS minute,
event_type,
countState() AS pv,
uniqExactState(user_id) AS uv,
uniqExactState(device_id) AS devices
FROM learn_ck.events_local
GROUP BY minute, event_type;
-- 6. 看板查询:用 *Merge 函数把状态合并成最终值
SELECT
minute,
event_type,
countMerge(pv) AS pv,
uniqExactMerge(uv) AS uv,
uniqExactMerge(devices) AS devices
FROM learn_ck.events_minute_agg
WHERE minute >= now() - INTERVAL 1 HOUR
GROUP BY minute, event_type
ORDER BY minute, event_type;写入 Kafka 一条消息,1-3 秒内 Grafana 看板的曲线就会动 —— 这就是「实时数仓」的本质。
8. 📌 与 Flink CDC / Spark 写入 CK 的对比
| 维度 | CK Kafka 引擎 | Flink CDC | Spark 批写 CK |
|---|---|---|---|
| 延迟 | 秒级(默认 7.5s 攒批) | 亚秒(流处理) | 分钟~小时(批) |
| 吞吐 | 单分区 100k+ msg/s | 集群线性扩展 | 极高(适合 TB 级回灌) |
| 复杂转换 | 仅能在 MV 写 SQL | Java/SQL 全能 | Spark SQL 全能 |
| 多源 JOIN | 只能在 CK 内 JOIN | 流式 JOIN 强 | 批 JOIN 强 |
| Exactly-once | At-least-once(重复用 ReplacingMergeTree 去重) | 端到端 EOS | EOS(基于 Checkpoint) |
| 运维复杂度 | 极低(CK 自带) | 高(Flink JM/TM) | 中 |
| 适用场景 | 简单 ETL / 业务直接落 CK | 复杂流处理 / 多源合流 | 历史回灌 / 离线 ETL |
经验建议:
- 95% 的「埋点 → CK」直接用 CK Kafka 引擎,越简单越好。
- 复杂的「多表 JOIN 后再落 CK」用 Flink CDC。
- 历史数据回灌(一次性 TB 级)用 Spark 或
clickhouse-client --input-format。
9. 本章小结
┌─────────────────────────────────────────────────────────────────┐
│ 本章核心要点 │
├─────────────────────────────────────────────────────────────────┤
│ │
│ ① Kafka 三件套 │
│ Kafka 表(消费者)+ MV(搬运工)+ MergeTree 表(存数据) │
│ 必记 SETTINGS:broker_list / topic_list / group_name / │
│ format / num_consumers / max_block_size │
│ │
│ ② MaterializedPostgreSQL │
│ 基于 PG 逻辑复制 slot 做 OLAP 镜像 │
│ 适合整库同步,不适合复杂转换 │
│ │
│ ③ 表函数家族 │
│ mysql() / postgresql() / s3() / hdfs() / remote() / │
│ clusterAllReplicas() —— 把 CK 当跨库联邦查询入口 │
│ │
│ ④ S3 / HDFS │
│ 可读 Parquet / CSV / ORC,可写归档 │
│ │
│ ⑤ ETL 工具栈选型 │
│ 简单:CK Kafka 引擎 │
│ CDC:Debezium → Kafka → CK │
│ 可视化:Airbyte │
│ 复杂流:Flink CDC │
│ │
│ ⑥ 实时数仓黄金链路 │
│ Kafka → 明细 → AggregatingMergeTree(MV 链)→ 看板 │
│ │
└─────────────────────────────────────────────────────────────────┘10. 面试高频题
Q1:ClickHouse 接 Kafka 为什么需要「三件套」(Kafka 表 + MV + MergeTree)?
考察点:理解 Kafka 引擎的「虚拟表 / 单次消费」本质。
标准答案:
- Kafka 引擎表本身不存数据:它是 ClickHouse 对 Kafka 的「连接抽象」,每次
SELECT等价于「从 Kafka 拉一批,消费完即扔」。所以你不能直接对它跑分析查询。 - 物化视图(MV)作为搬运工:CK 的 MV 是 INSERT 触发器 —— 每次 Kafka 表「冒出」新数据(CK 后台周期性 poll),MV 自动把 SELECT 结果 INSERT 到目标表。
- MergeTree 表做最终存储:所有真实查询都打这张表,享受列存 + 稀疏索引的所有性能优势。
- 解耦的好处:MV 可以做字段过滤、转换、补字段(如加
inserted_at);可以挂多个 MV 同时把数据分流到不同结构的目标表(明细一份、聚合一份、错误一份)。
加分项:
- 能讲清
kafka_num_consumers、kafka_max_block_size、kafka_handle_error_mode='stream'等关键参数。 - 能讲
ReplacingMergeTree+ 业务主键 ORDER BY 解决 at-least-once 的重复问题。 - 能提到一个常见坑:在 Kafka 表上跑 SELECT 会真的消费数据,调试时一定别手抖。
易错点:
- 别说「Kafka 表存数据」—— 它不存。
- 别在 Kafka 表上加 ORDER BY、PARTITION BY —— 它不是 MergeTree 引擎。
Q2:MaterializedPostgreSQL 引擎的工作原理是什么?什么场景适用?
考察点:理解 PG 逻辑复制 + CK 实时同步的机制。
标准答案:
- PG 端:开启
wal_level = logical后,PG 的 WAL 可以被「逻辑解码」成行级 change events(INSERT/UPDATE/DELETE 而不是物理页变更)。每个订阅者注册一个 replication slot,PG 会保证不删除该 slot 还没消费的 WAL。 - CK 端:
MaterializedPostgreSQL引擎用 PG 的复制协议(pgoutput plugin)连接 PG,注册一个 slot,周期性拉取变更。每张同步过来的表底层用ReplacingMergeTree按 PG 主键存储 + 去重。 - 适用场景:业务库(PG)做交易,分析库(CK 镜像)跑 OLAP,做读写分离。
- 不适合:复杂 ETL(应该用 Debezium + Kafka);高频 DDL(CK 不会自动跟随);多源合并。
加分项:
- 能提到「slot 不释放会让 PG WAL 撑爆磁盘」这个生产坑 —— CK 挂了一定要去 PG 上
pg_drop_replication_slot()。 - 能区分「逻辑复制 slot」和「物理流复制 slot」。
易错点:
- 别把它说成「批量同步工具」—— 它是实时同步。
- 别说「PG 的所有类型都能映射」—— 部分类型(如 PG 的 box / circle)会变成 String 或不支持。
Q3:mysql() 表函数和 MaterializedMySQL 引擎有什么区别?什么时候各用哪个?
考察点:能不能区分「ad-hoc 查询」和「持续同步」两种集成模式。
标准答案:
mysql()表函数:每次查询时实时连 MySQL 拉数据。CK 会把 WHERE 条件下推(predicate pushdown)到 MySQL 执行,但本质上每次查询都打 MySQL。MaterializedMySQL引擎:基于 MySQL binlog 的实时同步(类似 MaterializedPostgreSQL)。CK 在内部存全量数据,查询不打 MySQL。MySQL引擎表(持久化的 mysql 表函数):建表时绑定到一个 MySQL 表上,每次查询仍走 MySQL,只是不用每次写连接信息。- 如何选:
- 小数据量、偶尔 ad-hoc 查询 →
mysql()表函数 - 业务数据大、需要 CK 性能、对 MySQL 主库压力小 →
MaterializedMySQL引擎 - 需要做 MySQL → CK 的复杂 ETL → Debezium / Flink CDC
- 小数据量、偶尔 ad-hoc 查询 →
加分项:
- 能提到
MaterializedMySQL在新版本一直是 experimental,生产慎用。 - 能提到
Dictionary替代方案:把 MySQL 表灌进字典(第 9 章),用作 JOIN 的右表。
易错点:
- 别说「
mysql()表函数会把全表拉过来」—— CK 有 predicate pushdown。 - 别在
mysql()表上跑大表 JOIN —— 流量会非常恐怖。
Q4:什么是 CDC?Debezium 的工作原理是什么?怎么和 ClickHouse 配合?
考察点:CDC 概念 + 主流工具栈。
标准答案:
- CDC = Change Data Capture:把数据库的「写操作流」(INSERT / UPDATE / DELETE)实时捕获出来,变成可消费的事件流。
- Debezium 工作原理:作为 Kafka Connect 插件运行,连接源数据库的 change log(MySQL binlog、PG WAL、MongoDB oplog 等),把变更解析成 JSON / Avro 事件,按表打到 Kafka topic。每行变更包含
op(c/u/d)、before、after三个字段。 - CK 端消费:用 Kafka 引擎订阅 Debezium 的 topic,物化视图把
after字段拆出来 INSERT 到ReplacingMergeTree(按业务主键去重);DELETE 事件用version列 + Replacing 实现「软删除 → 自动覆盖」。 - 链路:MySQL → Debezium Connector → Kafka → CK Kafka 引擎 → MV → ReplacingMergeTree。
加分项:
- 能提到 Debezium 的 snapshot 模式:第一次启动时先全量快照再追 binlog。
- 能讲
__deleted字段如何处理删除:用is_deleted = 1列,配合ReplacingMergeTree(version, is_deleted)做 row-level delete 的最终一致。 - 能比较 Debezium 和 Flink CDC:Flink CDC 内置 Debezium,把 Source / Sink / 转换都做完了,省去 Kafka 中转。
易错点:
- 别说 Debezium 是数据库 —— 它是 Kafka Connect 插件。
- 别忽略「保序」问题:Debezium 默认按表分 partition,保证单行的更新有序。
Q5:用 s3() 表函数读 Parquet 文件有什么好处和限制?
考察点:数据湖与 OLAP 引擎的协作。
标准答案:
- 好处:
- 零拷贝接入:Spark / Flink 输出到 S3 的 Parquet 不需要再灌进 CK 就能查。
- 冷热分层:CK 内只存最近 30 天,30 天前归档到 S3,用户查老数据时透明打到 S3。
- 多文件并行读:CK 自动并行读 glob 匹配的多个 Parquet 文件。
- WHERE 下推:Parquet 自带 row group 统计信息,CK 能跳过不命中范围的 row group(类似稀疏索引)。
- 限制:
- 比本地 MergeTree 慢得多:网络 IO + Parquet 解码开销。
- 没有 CK 的稀疏主键索引和 Skip Index。
- 写入的 Parquet 文件每次都是新文件,不能 update / delete。
- 小文件 hell:太多小 Parquet 反而拖慢查询,最好用 Spark / DuckDB /
clickhouse-local提前合并。
加分项:
- 能提到
s3Cluster()表函数:在 CK 集群里并行读 S3,每个节点拉一部分文件。 - 能提到
INSERT INTO FUNCTION s3(...)把 CK 数据归档到 S3 是常见用法。 - 能讲「CK 也支持 IAM Role / EC2 Instance Profile,无需明文 AK/SK」。
易错点:
- 别用 s3() 当作热查询源 —— 性能不够。
- 别在数据快变大的场景下让 S3 文件无限增长,需要做 compaction。
Q6:remote() 和 Distributed 引擎是什么关系?
考察点:两种「跨节点查询」机制的差异。
标准答案:
Distributed引擎:在配置文件remote_servers里定义好集群拓扑,建一张持久化的「路由表」。读 / 写都走这张表,CK 自动 fan-out / fan-in。remote()表函数:临时指向某个远程节点 / 集群,不需要预先在remote_servers配集群。每次查询都现填连接信息(host、port、user、pwd)。- 配套的还有:
remoteSecure()—— TLS 版本cluster(<cluster_name>, db, table)—— 在已配置的集群里跑(自动 fan-out)clusterAllReplicas(...)—— 强制查询每个分片的所有副本(用于一致性比对)
- 何时各用哪个:
- 稳定业务读写 →
Distributed表 - 临时数据搬家、跨集群 ad-hoc 查询 →
remote() - 集群已配好,临时跑分布式查询 →
cluster()/clusterAllReplicas()
- 稳定业务读写 →
加分项:
- 能给经典用法
INSERT INTO new SELECT * FROM remote('old', ...)做集群迁移。 - 能讲性能:
remote()现填鉴权信息,每次查询都要重连,频繁查询不如Distributed。
易错点:
- 别说
remote()只能查单点 —— 第一个参数也可以是集群名(按集群配置 fan-out)。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
"""
第 14 章 · 集成生态 · integration_play.py
演示内容:
1. mysql() 表函数:跨库查询
2. postgresql() 表函数:跨库查询
3. s3() 表函数:读公开 Parquet(无需 AK/SK 的公开桶)
4. url() 表函数:读公开 CSV
5. remote() 表函数:跨集群 ad-hoc
凭证 / 地址全部从【环境变量】读取,避免硬编码:
CK_HOST = 127.0.0.1 (CK HTTP 端口 8123)
MYSQL_HOST / DB / USER / PWD / TABLE
PG_HOST / DB / USER / PWD / TABLE / SCHEMA
S3_URL = https://<bucket>/path/*.parquet
S3_AK / S3_SK = (可选)
REMOTE_CK = other-ck:9000
REMOTE_USER / REMOTE_PWD
不设置环境变量也能跑:会自动跳过对应章节,仅演示 ClickHouse 端的 SQL 模板。
依赖:
pip install clickhouse-connect
"""
from __future__ import annotations
import os
import sys
from contextlib import contextmanager
try:
import clickhouse_connect
except ImportError:
print("请先 pip install clickhouse-connect")
sys.exit(1)
CK_HOST = os.getenv("CK_HOST", "127.0.0.1")
CK_PORT = int(os.getenv("CK_PORT", "8123"))
CK_USER = os.getenv("CK_USER", "default")
CK_PWD = os.getenv("CK_PWD", "")
CK_DB = os.getenv("CK_DB", "learn_ck")
def banner(title: str) -> None:
print("\n" + "=" * 78)
print(f" {title}")
print("=" * 78)
@contextmanager
def section(title: str):
banner(title)
try:
yield
except Exception as exc: # noqa: BLE001
print(f" ⚠ 跳过:{type(exc).__name__}: {exc}")
def main() -> None:
print(f"连接 {CK_HOST}:{CK_PORT}(user={CK_USER})...")
client = clickhouse_connect.get_client(
host=CK_HOST, port=CK_PORT, username=CK_USER, password=CK_PWD,
database="default",
)
version = client.query("SELECT version()").result_rows[0][0]
print(f"OK · ClickHouse 版本 {version}")
client.command(f"CREATE DATABASE IF NOT EXISTS {CK_DB}")
# =====================================================================
# 1) mysql() 表函数
# =====================================================================
with section("1) mysql() 表函数(跨库查询 MySQL)"):
host = os.getenv("MYSQL_HOST")
if not host:
print(" 未设置 MYSQL_HOST 环境变量,仅打印 SQL 模板:")
print("""
SELECT user_id, count() AS cnt
FROM mysql('mysql-host:3306','db','table','user','pwd')
WHERE created_at >= today() - 7
GROUP BY user_id
ORDER BY cnt DESC
LIMIT 5;
""".strip())
else:
db = os.getenv("MYSQL_DB", "shop")
tbl = os.getenv("MYSQL_TABLE", "users")
user = os.getenv("MYSQL_USER", "root")
pwd = os.getenv("MYSQL_PWD", "")
sql = f"""
SELECT count() FROM mysql('{host}', '{db}', '{tbl}', '{user}', '{pwd}')
"""
res = client.query(sql).result_rows
print(f" MySQL {host}/{db}.{tbl} 行数 = {res[0][0]}")
# =====================================================================
# 2) postgresql() 表函数
# =====================================================================
with section("2) postgresql() 表函数(跨库查询 PostgreSQL)"):
host = os.getenv("PG_HOST")
if not host:
print(" 未设置 PG_HOST 环境变量,仅打印 SQL 模板:")
print("""
SELECT * FROM postgresql(
'pg-host:5432','db','table','user','pwd','public'
) WHERE status = 'paid' LIMIT 10;
""".strip())
else:
db = os.getenv("PG_DB", "shop")
tbl = os.getenv("PG_TABLE", "orders")
user = os.getenv("PG_USER", "postgres")
pwd = os.getenv("PG_PWD", "")
schema = os.getenv("PG_SCHEMA", "public")
sql = f"""
SELECT count() FROM postgresql(
'{host}','{db}','{tbl}','{user}','{pwd}','{schema}'
)
"""
res = client.query(sql).result_rows
print(f" PG {host}/{db}.{schema}.{tbl} 行数 = {res[0][0]}")
# =====================================================================
# 3) url() 表函数:读公开 CSV
# =====================================================================
with section("3) url() 表函数:把网络上的 CSV 当成 CK 表查"):
# 用 CK 官方公开数据集做演示(如果环境能联网)
url = os.getenv(
"DEMO_URL",
"https://datasets.clickhouse.com/hits/tsv/hits_v1.tsv.gz"
)
# 仅打印 SQL,不真去 download(百 MB 级别)
print(f" 示例 SQL(不实际执行重型下载):")
print(f"""
SELECT count() FROM url(
'{url}',
'TSV',
'WatchID UInt64, ...'
);
""".strip())
# 真正跑一个轻量级的 demo:读一个非常小的 CSV
small_url = os.getenv(
"SMALL_CSV",
"https://gist.githubusercontent.com/curran/a08a1080b88344b0c8a7/raw/iris.csv"
)
try:
sql = f"""
SELECT count() FROM url(
'{small_url}', 'CSVWithNames'
)
"""
res = client.query(sql).result_rows
print(f" 小 CSV 行数 = {res[0][0]} (URL: {small_url})")
except Exception as exc: # noqa: BLE001
print(f" ⚠ 网络受限,跳过:{exc}")
# =====================================================================
# 4) s3() 表函数
# =====================================================================
with section("4) s3() 表函数:直接读 S3 上的 Parquet/CSV"):
s3_url = os.getenv("S3_URL")
if not s3_url:
print(" 未设置 S3_URL 环境变量,仅打印 SQL 模板:")
print("""
-- 读
SELECT count() FROM s3(
'https://my-bucket.s3.amazonaws.com/events/*.parquet',
'Parquet',
'AKIA...', 'SECRET...'
);
-- 写(归档)
INSERT INTO FUNCTION s3(
'https://my-bucket.s3.amazonaws.com/archive/2026.parquet',
'Parquet'
)
SELECT * FROM learn_ck.events_local
WHERE event_time >= '2026-01-01';
""".strip())
else:
ak = os.getenv("S3_AK", "")
sk = os.getenv("S3_SK", "")
fmt = os.getenv("S3_FORMAT", "Parquet")
if ak and sk:
sql = f"SELECT count() FROM s3('{s3_url}', '{fmt}', '{ak}', '{sk}')"
else:
sql = f"SELECT count() FROM s3('{s3_url}', '{fmt}')"
res = client.query(sql).result_rows
print(f" S3 {s3_url} 行数 = {res[0][0]}")
# =====================================================================
# 5) remote() 表函数
# =====================================================================
with section("5) remote() 表函数:跨集群 / 跨节点 ad-hoc"):
rh = os.getenv("REMOTE_CK")
if not rh:
print(" 未设置 REMOTE_CK 环境变量,仅打印 SQL 模板:")
print("""
SELECT count() FROM remote(
'other-ck:9000', learn_ck.events_local, 'user', 'pwd'
);
-- 跨集群迁移:
INSERT INTO new_db.events
SELECT * FROM remote('old_cluster_name', old_db.events);
""".strip())
else:
ru = os.getenv("REMOTE_USER", "default")
rp = os.getenv("REMOTE_PWD", "")
res = client.query(
f"SELECT count() FROM remote('{rh}', system.tables, '{ru}', '{rp}')"
).result_rows
print(f" 远端 system.tables 行数 = {res[0][0]}")
print("\n演示完成。")
if __name__ == "__main__":
main()