一、先弄明白这条实时同步链路

MySQL 和 ClickHouse 就像两个性格完全不同的人。MySQL 是精打细算的小会计,适合把每一笔账记清楚,修改、删除都特别方便。ClickHouse 是雷厉风行的数据分析师,你问它“过去一年每个用户总共花了多少钱”,它几秒钟就能算出来,但如果你让它频繁地删除或修改某个订单,它反而会很不乐意。

正因为两者各有脾气,很多团队的做法是:让 MySQL 继续扮演“业务主力”,专门处理订单、用户、支付这些实时事务;再让 ClickHouse 接一份数据过去,用来做复杂的统计和报表。可是问题来了:MySQL 里的数据变了,ClickHouse 那边的数据怎么跟着变?总不能每次都用人工搬吧。

于是,一条实时同步链路出现了:Canal 负责盯着 MySQL 的每一次变化,然后把变化“打包”传走;ClickHouse 那边用物化视图自动接收并整理这些变化。听起来很完美,实际用起来却发现,Canal 和物化视图常常会“打架”。这篇文章就是想把打架的原因说清楚,再告诉你怎样让它们配合默契。

1.1 实时同步到底要解决什么

说白了,就是一件事:MySQL 数据变了,ClickHouse 也要跟着变,而且最好“秒级”变。比如你在 MySQL 里新建了一笔订单,ClickHouse 的每日销售统计里最好立刻就能体现出来,而不是等到第二天跑批。

要实现这个效果,首先得能“捕获”MySQL 的变化。MySQL 里有个东西叫 binlog,它把每一笔插入、修改、删除都记录在案。Canal 就是专门读 binlog 的小能手。另外一步就是 ClickHouse 这边得能“消费”这些变化,把变化结果写进自己的表里。物化视图正好能帮上忙,它像是一段自动触发的程序,一旦有数据进来,就按照你写好的规则去处理。

1.2 一个容易让人忽略的痛点

很多开发者第一次搭这条链路时,会想:Canal 把 MySQL 的数据原封不动地丢给 ClickHouse,然后 ClickHouse 用物化视图做个统计,不就完事了吗?现实没这么简单。因为 MySQL 是支持任意更新和删除的,而 ClickHouse 对更新删除的支持却很“拧巴”。这中间的差异,就是冲突的老家。

二、Canal:专门“盯梢” MySQL 的哨兵

Canal 是一个开源工具,你的 MySQL 只要开启 binlog,Canal 就能把自己伪装成一个“从库”,然后持续不断读取 binlog 的变更记录。说得通俗一点,它就像蹲在 MySQL 门口的哨兵,每进来一笔数据,它就记一下;哪儿改了一笔,它也记一下;哪儿删了一笔,它也记一下。

它的输出,通常是一串 JSON 格式的消息。每条消息里包含表名、操作类型、变更前后的字段值等等。下游系统拿到这份消息,就可以知道 MySQL 里发生了什么。

2.1 启动 Canal 的示例

下面我们用 Docker 启动一个 Canal 容器,让它监听一台 MySQL 实例。技术栈统一使用 Shell 命令来演示。

# 技术栈:Shell
# 注意:所有命令都在命令行里执行,这里只是演示核心步骤

# 第一步:启动一个开启 binlog 的 MySQL 容器(简化版)
docker run -d --name mysql-demo \
  -e MYSQL_ROOT_PASSWORD=root \
  mysql:8.0 \
  --binlog-format=ROW \
  --server-id=1

# 第二步:启动 Canal,并让它连接刚才的 MySQL
# canal.instance.master.address 表示 MySQL 的地址和端口
# canal.instance.dbUsername/dbPassword 是 MySQL 账号密码
docker run -d --name canal-demo \
  --link mysql-demo \
  -p 11111:11111 \
  -e canal.instance.master.address=mysql-demo:3306 \
  -e canal.instance.dbUsername=root \
  -e canal.instance.dbPassword=root \
  -e canal.instance.connectionCharset=UTF-8 \
  -e canal.instance.defaultDatabaseName=test \
  -e canal.mq.topic=order-binlog \
  canal/canal-server

启动后,Canal 就会自动连接 MySQL,开始监听。你可以在 MySQL 里执行一条 INSERTUPDATEDELETE,Canal 都会把它转成一条变更消息。这些消息通常会被发送到消息队列,比如 Kafka,也可以直接用 TCP 方式输出。

三、物化视图:ClickHouse 里的“自动搬运工”

ClickHouse 的物化视图,英文叫 Materialized View。它的作用听起来也很朴素:当你往某张表里插入新的数据时,它会自动触发一段查询,再把查询结果写到另一张表里。为了避免晦涩,你可以把它想象成一个“自动分拣机”——传送带上如果出现新包裹,它就会自动按标签分类,放到对应的柜子里。

3.1 一个最简单的物化视图

假设 ClickHouse 里有一张原始订单消息表,只要这张表里来了一条新记录,我们希望自动把它写进另一张“订单详情表”。用 SQL 表达就是创建一个从源表到目标表的物化视图。这里我们用 Shell 调用 clickhouse-client 来执行 SQL。

# 技术栈:Shell
# 在 ClickHouse 中执行一段查询命令

# 先建一张源表:用来接收 Canal 发来的原始数据(简化字段)
clickhouse-client --query "
CREATE TABLE IF NOT EXISTS app.incoming_orders (
    order_id     UInt64,
    user_id      UInt64,
    amount       Decimal(10,2),
    event_time   DateTime
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(event_time)
ORDER BY (order_id)
"

# 再建一张目标表:专门存放处理后的明细
clickhouse-client --query "
CREATE TABLE IF NOT EXISTS app.processed_orders (
    order_id     UInt64,
    user_id      UInt64,
    amount       Decimal(10,2),
    event_time   DateTime
) ENGINE = MergeTree()
ORDER BY order_id
"

# 创建物化视图:incoming_orders 一旦有新数据,就自动复制到 processed_orders
clickhouse-client --query "
CREATE MATERIALIZED VIEW app.order_mv TO app.processed_orders
AS SELECT order_id, user_id, amount, event_time
FROM app.incoming_orders
"

写完这段之后,只要有人往 app.incoming_orders 里插数据,app.processed_orders 里几乎同时会出现一样的数据。你不需要手动去写 INSERT INTO processed_orders SELECT ...,物化视图会帮你搞定。

3.2 物化视图适合干什么

物化视图特别适合做“追加型”数据流的处理。比如:Canal 把 MySQL 的新增订单一条条发过来,物化视图自动按用户分组做汇总,然后持续更新一张统计表。注意关键词“新增”。如果 MySQL 里只是不断新增数据,不修改、不删除,那么物化视图几乎完美。

四、冲突:为什么两者会“打架”

既然 Canal 负责把数据送进来,物化视图负责自动处理,为什么还会打架呢?核心原因有三点。

4.1 重复数据让人头大

Canal 在运行过程中很可能会重复投递消息。比如网络抖动,Canal 重启,或者消费端没及时确认已读,都会导致同一条 binlog 变更被发送两次。如果你的流程是“Canal → ClickHouse 源表 → 物化视图”,那么源表里会出现两条一模一样的记录。物化视图可不会帮你识别重复,它会照样把两条记录都算进去。结果就是:订单金额被统计了两次,报表数据直接翻倍。

这就像有人往自动分拣机上塞了两件一模一样的包裹,分拣机只会机械地分两次,不会说“哎,这个刚才已经来过了”。

4.2 MySQL 的更新和删除是“硬伤”

MySQL 里很常见的 UPDATE 操作,对 Canal 来说是一条“修改”事件。可 ClickHouse 的物化视图只认“插入”。你无法简单地对 ClickHouse 说“把某一行改成新值”。就算你把 MySQL 的修改操作变成“先删后插”,ClickHouse 也不能真正删除旧数据,它只能逻辑删除或者让旧数据放在那里。

如果物化视图直接统计源表,MySQL 里把订单金额从 100 改成 200,ClickHouse 这边会看到两条记录:一条 100,一条 200。用户最终统计出来的金额会变成 300,而不是 200。这不就乱套了吗?

4.3 上游延迟带来的“假实时”

Canal 本身是实时的,但网络传输、消息队列积压、ClickHouse 写入瓶颈,都可能让数据晚到。物化视图只负责处理“已经进入源表”的数据,它不会知道 MySQL 里还有多少数据正在路上。于是你会看到 ClickHouse 里的报表一会儿缺数据,一会儿又突然跳出一大堆补发数据。这就是所谓的“假实时”。

五、协同:怎么让它们同心协力

既然矛盾集中在“更新、删除、重复”上,那我们就得想办法把这三件事绕过去。聪明的工程师们想出了一个思路:让 Canal 发到 Kafka,让 ClickHouse 通过“Kafka 引擎表”接入数据,再让物化视图把 Kafka 表中的数据“搬运”到目标表。同时,利用 ClickHouse 的 ReplacingMergeTree 来去重,解决 MySQL 的更新与删除问题。

5.1 引入 Kafka 作为缓冲

Kafka 是一个消息队列,它像一个巨大的“邮件中转站”。Canal 负责把 MySQL 的变更消息扔进 Kafka,ClickHouse 负责从 Kafka 里取消息。好处是:即使 ClickHouse 暂时宕机,消息也不会丢;同时可以重复消费,方便修复数据。

架构变成这样:

Canal → Kafka → ClickHouse 的 Kafka 引擎表 → 物化视图 → 真正的存储表

5.2 创建 ClickHouse 的 Kafka 引擎表

Kafka 引擎表专门用于连接 Kafka 主题。它本身不存数据,只是“管道”的入口。我们用 Shell 命令创建一个 Kafka 引擎表。

# 技术栈:Shell
# 创建一张 Kafka 引擎表,用来接收 order-binlog 主题的消息
clickhouse-client --query "
CREATE TABLE IF NOT EXISTS app.kafka_order_events (
    order_id   UInt64,
    user_id    UInt64,
    amount     Decimal(10,2),
    event_time DateTime,
    version    UInt64,        -- 版本号,可以是 binlog 的 position
    operation  String         -- 操作类型:insert / update / delete
) ENGINE = Kafka()
SETTINGS
    kafka_broker_list = 'kafka:9092',
    kafka_topic_list = 'order-binlog',
    kafka_group_name = 'clickhouse-consumer-group',
    kafka_format = 'JSONEachRow',
    kafka_num_consumers = 1
"

这张表创建之后,只要 Kafka 里有消息,就会源源不断进入这张“虚拟表”。

5.3 用物化视图把 Kafka 表数据写入存储表

我们接着创建一张真正的存储表,并让物化视图把 Kafka 引擎表里的数据写进去。注意,这张存储表要用 ReplacingMergeTree(version) 引擎,这样当同一条 order_id 有了新版本时,旧版本会在后台被替换掉。

# 技术栈:Shell
# 创建真正的存储表,使用 ReplacingMergeTree,按 order_id 去重,取 version 最大的一条
clickhouse-client --query "
CREATE TABLE IF NOT EXISTS app.order_events_final (
    order_id   UInt64,
    user_id    UInt64,
    amount     Decimal(10,2),
    event_time DateTime,
    version    UInt64,
    operation  String
) ENGINE = ReplacingMergeTree(version)
ORDER BY order_id
"

# 创建物化视图:把 Kafka 引擎表中的每一条消息都写入 order_events_final
clickhouse-client --query "
CREATE MATERIALIZED VIEW app.order_events_mv TO app.order_events_final
AS SELECT
    order_id,
    user_id,
    amount,
    event_time,
    version,
    operation
FROM app.kafka_order_events
"

完成之后,Canal 发来一条 Kafka 消息,物化视图就会把这条消息插入 order_events_final。如果同一条 order_id 来了两遍,ReplacingMergeTree 会在后台根据 version 保留最新的那一条。这样,MySQL 里的更新操作就能被模拟出来了。

5.4 小心“物化视图”和“替换引擎”的配合陷阱

这里必须提醒你一个关键点:物化视图是“来一条处理一条”,它不会在插入的时候去重。也就是说,即使 order_events_final 使用了 ReplacingMergeTree,物化视图也依然会把两条重复数据都写进去。ReplacingMergeTree 的去重行为发生在后台的“合并”阶段,而不是数据到达的那一秒。

所以在查询时,如果你想立刻得到精确结果,必须要用 FINAL 关键字,或者通过 GROUP BY 自己兜底。例如:

# 技术栈:Shell
# 查询时使用 FINAL,强制结果只保留每个 order_id 的最新版本
clickhouse-client --query "
SELECT
    order_id,
    argMax(user_id, version)    AS user_id,
    argMax(amount, version)     AS amount,
    argMax(event_time, version) AS event_time,
    argMax(operation, version)  AS operation
FROM app.order_events_final
GROUP BY order_id
"

如果你既想用物化视图做自动汇总,又想保证结果不重复,那你最好让物化视图把数据“原样”写到这张按 order_id 去重的表里,然后再从这张表做统计查询,而不是直接让物化视图做 sum 聚合。

5.5 处理 MySQL 的删除操作

MySQL 里删除一行,Canal 会发出一条 operation = "delete" 的消息。但 ClickHouse 里没有真正的删除,所以你的同步程序(或者物化视图)可以把这条消息变成“更新版本号”,同时把某个字段标记为已删除。比如增加一个 is_deleted 字段,当 operation = 'delete' 时,写入一条 is_deleted = 1 的记录。查询时,你在 FINAL 结果中过滤掉 is_deleted = 1 即可。

示例如下:

# 技术栈:Shell
# 在存储表中增加 is_deleted 字段,把 MySQL 的删除变成一条新版本记录
clickhouse-client --query "
CREATE TABLE IF NOT EXISTS app.order_events_final (
    order_id    UInt64,
    user_id     UInt64,
    amount      Decimal(10,2),
    event_time  DateTime,
    version     UInt64,
    operation   String,
    is_deleted  UInt8 DEFAULT 0
) ENGINE = ReplacingMergeTree(version)
ORDER BY order_id
"

# 查询时排除已删除的数据
clickhouse-client --query "
SELECT
    order_id,
    argMax(user_id, version)    AS user_id,
    argMax(amount, version)     AS amount,
    argMax(event_time, version) AS event_time
FROM app.order_events_final
WHERE argMax(is_deleted, version) = 0
GROUP BY order_id
"

这算是目前比较可行的“伪删除”方案。

六、应用场景和优缺点

6.1 适合的场景

这种“Canal + 物化视图”的架构,最适合那些以“新增为主、修改很少”的业务。比如:

  • 订单流水记录:订单一旦创建,很少改金额。
  • 用户行为日志:用户点击、浏览都是只增不改。
  • 监控告警事件:报警事件发生后,状态偶尔变一下,但不会频繁改。

在这些场景下,物化视图的简单直接就能发挥巨大优势。

6.2 优点

  • 同步链路简单,不引入太多额外计算引擎。
  • 物化视图自动触发,代码量少,维护成本低。
  • Canal 对 MySQL 无侵入,只是读取 binlog。
  • ClickHouse 秒级聚合能力强,报表响应快。

6.3 缺点

  • 无法直接处理 MySQL 的更新与删除,需要做不少“纸面功夫”。
  • 物化视图不会自动去重,需要配合 ReplacingMergeTree + FINAL,查询时会有性能损耗。
  • 链路中有 Kafka、Canal、ClickHouse 多个组件,任何一个出问题都会导致数据延迟。
  • 对数据一致性要求极高的场景(比如财务对账),这套方案得再加校验机制。

七、注意事项

  1. binlog 必须开启 ROW 格式。Canal 只有读 ROW 格式才能拿到变更前后的数据,否则无法工作。
  2. Canal 的分区与稳定性。Canal 服务可以多实例部署,但要注意同一个 binlog 位点只能被一个实例消费,否则会重复。
  3. Kafka 的消费者组。ClickHouse 的 Kafka 引擎表默认用自己的消费者组,如果你重放消息,记得调整 consumer group 或重置 offset。
  4. 物化视图的幂等性。物化视图不是任务调度器,它不会重试,也不保证消息顺序。如果 Kafka 乱序,你的版本号机制要能兜住乱序情况。
  5. ClickHouse 版本选择。不同版本的 ClickHouse 对 Kafka 引擎表、物化视图的行为略有差异,生产环境前一定要在测试环境验证。
  6. 数据校验。每隔一段时间,要抽查 MySQL 里的关键数据与 ClickHouse 里的一致性,防止长期积累出现偏差。

八、总结

Canal 和物化视图本身并不冲突,它们只是在面对 MySQL 的“更新和删除”时显得力不从心。只要我们把 MySQL 的变更事件抽象成“带版本号的追加消息”,再通过 Kafka 传输给 ClickHouse,最后用物化视图配合 ReplacingMergeTree 去重,就能让整个链路跑得又稳又准。

这套架构的精髓在于:物化视图负责“快”,Canal 负责“准”,而 Kafka 负责“缓冲”。快、准、缓冲,三者协同,才能应对真实业务里的数据变化。

希望这篇文章能帮你少踩几个坑。如果下次有人问你 MySQL 到 ClickHouse 怎么实时同步,你可以告诉他:别让 Canal 和物化视图硬碰硬,给它们中间加一层 Kafka,再给 ClickHouse 加一张能去重的表,事情就顺了。