Skip to content

第 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 / PulsarKafka 引擎 + MV实时摄入埋点、日志、CDC
MySQLMySQL 引擎 / mysql() 表函数跨库 join、维表
PostgreSQLPostgreSQL 引擎 / postgresql() 表函数同上
PostgreSQL(实时同步)MaterializedPostgreSQL 引擎整库做 OLAP 镜像
S3 / OSS / MinIOS3 引擎 / s3() 表函数数据湖直读 Parquet
HDFSHDFS 引擎 / hdfs() 表函数Hadoop 生态联邦
任意 CK 集群remote() / clusterAllReplicas()跨集群联邦查询
通用 ETLAirbyte / 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_listKafka 集群地址至少 3 个 broker
kafka_topic_list订阅的 topic(可多个,逗号分隔)业务对应
kafka_group_nameKafka 消费组名每个独立消费链路一个
kafka_format消息格式JSONEachRow / CSV / Avro / Protobuf
kafka_num_consumers同一节点的消费者线程数≤ topic 的分区数
kafka_max_block_size一次拉取多少条凑成一个 Part65536(默认 8192 偏小)
kafka_skip_broken_messages跳过多少条解析失败的消息100(生产可改更大或 0 死循环)
kafka_thread_per_consumer是否每个消费者一个线程1(默认 0)
kafka_commit_every_batch每批 commit offset0(默认按时间 commit)
kafka_handle_error_mode错误处理:default / streamstream 把错误也送到 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.errorstext_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 注意事项

  1. 每张同步过来的表底层引擎是 ReplacingMergeTree,按 PG 主键去重。
  2. DDL 变更需要重建:CK 不会自动跟随 PG 的 ALTER TABLE
  3. WAL slot 不释放会撑爆 PG 磁盘:CK 挂了一定要去 PG 上 pg_drop_replication_slot() 清理。
  4. 不支持 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() 所有副本都查(看一致性)

最常见的两个用途:

  1. 数据搬家INSERT INTO new_cluster.events_all SELECT * FROM remote('old_cluster', ...)
  2. 数据回查 / 比对:跨集群对账,看新老集群结果是否一致。

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 引擎 ─→ MergeTree

CK 端:用 Kafka 引擎消费 Debezium 输出的格式(JSONEachRowAvro),用 ReplacingMergeTree 按业务主键去重 + 用 __deleted 列处理删除事件。

6.2 Airbyte:低代码 EL 框架

Airbyte 提供 300+ Source / Destination 连接器,UI 拖拽就能配 MySQL → ClickHouse / Salesforce → ClickHouse 这种链路。适合数据团队人手不足、又要接很多业务源的场景。

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 看板的曲线就会动 —— 这就是「实时数仓」的本质。


维度CK Kafka 引擎Flink CDCSpark 批写 CK
延迟秒级(默认 7.5s 攒批)亚秒(流处理)分钟~小时(批)
吞吐单分区 100k+ msg/s集群线性扩展极高(适合 TB 级回灌)
复杂转换仅能在 MV 写 SQLJava/SQL 全能Spark SQL 全能
多源 JOIN只能在 CK 内 JOIN流式 JOIN 强批 JOIN 强
Exactly-onceAt-least-once(重复用 ReplacingMergeTree 去重)端到端 EOSEOS(基于 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 引擎的「虚拟表 / 单次消费」本质。

标准答案

  1. Kafka 引擎表本身不存数据:它是 ClickHouse 对 Kafka 的「连接抽象」,每次 SELECT 等价于「从 Kafka 拉一批,消费完即扔」。所以你不能直接对它跑分析查询。
  2. 物化视图(MV)作为搬运工:CK 的 MV 是 INSERT 触发器 —— 每次 Kafka 表「冒出」新数据(CK 后台周期性 poll),MV 自动把 SELECT 结果 INSERT 到目标表。
  3. MergeTree 表做最终存储:所有真实查询都打这张表,享受列存 + 稀疏索引的所有性能优势。
  4. 解耦的好处:MV 可以做字段过滤、转换、补字段(如加 inserted_at);可以挂多个 MV 同时把数据分流到不同结构的目标表(明细一份、聚合一份、错误一份)。

加分项

  • 能讲清 kafka_num_consumerskafka_max_block_sizekafka_handle_error_mode='stream' 等关键参数。
  • 能讲 ReplacingMergeTree + 业务主键 ORDER BY 解决 at-least-once 的重复问题。
  • 能提到一个常见坑:在 Kafka 表上跑 SELECT 会真的消费数据,调试时一定别手抖。

易错点

  • 别说「Kafka 表存数据」—— 它不存。
  • 别在 Kafka 表上加 ORDER BY、PARTITION BY —— 它不是 MergeTree 引擎。

Q2:MaterializedPostgreSQL 引擎的工作原理是什么?什么场景适用?

考察点:理解 PG 逻辑复制 + CK 实时同步的机制。

标准答案

  1. PG 端:开启 wal_level = logical 后,PG 的 WAL 可以被「逻辑解码」成行级 change events(INSERT/UPDATE/DELETE 而不是物理页变更)。每个订阅者注册一个 replication slot,PG 会保证不删除该 slot 还没消费的 WAL。
  2. CK 端MaterializedPostgreSQL 引擎用 PG 的复制协议(pgoutput plugin)连接 PG,注册一个 slot,周期性拉取变更。每张同步过来的表底层用 ReplacingMergeTree 按 PG 主键存储 + 去重。
  3. 适用场景:业务库(PG)做交易,分析库(CK 镜像)跑 OLAP,做读写分离。
  4. 不适合:复杂 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 查询」和「持续同步」两种集成模式。

标准答案

  1. mysql() 表函数:每次查询时实时连 MySQL 拉数据。CK 会把 WHERE 条件下推(predicate pushdown)到 MySQL 执行,但本质上每次查询都打 MySQL。
  2. MaterializedMySQL 引擎:基于 MySQL binlog 的实时同步(类似 MaterializedPostgreSQL)。CK 在内部存全量数据,查询不打 MySQL。
  3. MySQL 引擎表(持久化的 mysql 表函数):建表时绑定到一个 MySQL 表上,每次查询仍走 MySQL,只是不用每次写连接信息。
  4. 如何选
    • 小数据量、偶尔 ad-hoc 查询mysql() 表函数
    • 业务数据大、需要 CK 性能、对 MySQL 主库压力小MaterializedMySQL 引擎
    • 需要做 MySQL → CK 的复杂 ETL → Debezium / Flink CDC

加分项

  • 能提到 MaterializedMySQL 在新版本一直是 experimental,生产慎用。
  • 能提到 Dictionary 替代方案:把 MySQL 表灌进字典(第 9 章),用作 JOIN 的右表。

易错点

  • 别说「mysql() 表函数会把全表拉过来」—— CK 有 predicate pushdown。
  • 别在 mysql() 表上跑大表 JOIN —— 流量会非常恐怖。

Q4:什么是 CDC?Debezium 的工作原理是什么?怎么和 ClickHouse 配合?

考察点:CDC 概念 + 主流工具栈。

标准答案

  1. CDC = Change Data Capture:把数据库的「写操作流」(INSERT / UPDATE / DELETE)实时捕获出来,变成可消费的事件流。
  2. Debezium 工作原理:作为 Kafka Connect 插件运行,连接源数据库的 change log(MySQL binlog、PG WAL、MongoDB oplog 等),把变更解析成 JSON / Avro 事件,按表打到 Kafka topic。每行变更包含 op(c/u/d)、beforeafter 三个字段。
  3. CK 端消费:用 Kafka 引擎订阅 Debezium 的 topic,物化视图把 after 字段拆出来 INSERT 到 ReplacingMergeTree(按业务主键去重);DELETE 事件用 version 列 + Replacing 实现「软删除 → 自动覆盖」。
  4. 链路: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 引擎的协作。

标准答案

  1. 好处
    • 零拷贝接入:Spark / Flink 输出到 S3 的 Parquet 不需要再灌进 CK 就能查。
    • 冷热分层:CK 内只存最近 30 天,30 天前归档到 S3,用户查老数据时透明打到 S3。
    • 多文件并行读:CK 自动并行读 glob 匹配的多个 Parquet 文件。
    • WHERE 下推:Parquet 自带 row group 统计信息,CK 能跳过不命中范围的 row group(类似稀疏索引)。
  2. 限制
    • 比本地 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 引擎是什么关系?

考察点:两种「跨节点查询」机制的差异。

标准答案

  1. Distributed 引擎:在配置文件 remote_servers 里定义好集群拓扑,建一张持久化的「路由表」。读 / 写都走这张表,CK 自动 fan-out / fan-in。
  2. remote() 表函数临时指向某个远程节点 / 集群,不需要预先在 remote_servers 配集群。每次查询都现填连接信息(host、port、user、pwd)。
  3. 配套的还有
    • remoteSecure() —— TLS 版本
    • cluster(<cluster_name>, db, table) —— 在已配置的集群里跑(自动 fan-out)
    • clusterAllReplicas(...) —— 强制查询每个分片的所有副本(用于一致性比对)
  4. 何时各用哪个
    • 稳定业务读写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()

integration_play.py ↗