首页
/ Pathway × Debezium:用 Change Data Capture 实现 MongoDB 实时流处理与聚合回写

Pathway × Debezium:用 Change Data Capture 实现 MongoDB 实时流处理与聚合回写

2026-09-05 20:51:53作者:卓艾滢Kingsley

本文以 Pathway 仓库中的 MongoDB + Debezium 官方教程及其配套示例项目为主线,讲解如何为 MongoDB 这类非流式数据库搭建完整的变更数据捕获(CDC)链路:MongoDB 复制集 → Debezium → Kafka/ZooKeeper → Pathway,再将实时聚合结果通过 pw.io.mongodb.write 回写到 MongoDB。读完本文,你将能够独立部署这条五容器架构,理解 pw.io.debezium.readpw.io.mongodb.write 两个连接器的完整参数,并结合仓库源码掌握其底层实现细节。

为什么 MongoDB 需要 Debezium 做 CDC

MongoDB 等传统数据库本身并不面向流式场景构建,无法直接对外推送数据变更。要在 Pathway Live Data Framework 中把 MongoDB 集合当作输入流,必须引入一个 CDC 机制持续监听数据库并产出变更事件。该教程的方案是:

  • 输入侧:使用 Debezium 捕获 MongoDB 变更,经由 Kafka 发布到 topic,再由 Pathway 的 Debezium 连接器 pw.io.debezium.read 消费;
  • 输出侧:使用原生 MongoDB 连接器 pw.io.mongodb.write 将处理结果写回 MongoDB。

如果目标库是 PostgreSQL,仓库中也提供结构类似、但配置细节不同的教程,可对照参考。需要说明的是:Debezium 消息本质上走的是 Kafka,Debezium 连接器与常规 Kafka 连接器的唯一区别在于期望的消息格式(Debezium 封装格式)。

一个具体场景:实时求和

教程用一个简单场景贯穿全文:MongoDB 中有一个集合 values,文档包含一个 value 字段;数据持续插入,我们需要实时计算这些值的总和,并把结果存到另一个集合 sum_values 中。核心流水线只有三行:

import pathway as pw

# 若要使用 Pathway Live Data Framework Scale 的高级特性,
# 可获取免费 license key 并粘贴在此;使用 Community 版则注释掉下面这行。
pw.set_license_key("demo-license-key-with-telemetry")

# Kafka 连接设置
input_rdkafka_settings = {
    "bootstrap.servers": "kafka:9092",
    "security.protocol": "plaintext",
    "group.id": "0",
    "session.timeout.ms": "6000",
    "auto.offset.reset": "earliest",
}


class InputSchema(pw.Schema):
    value: int


if __name__ == "__main__":
    # 监听 Debezium 发布在 "my_mongo_db.test_database.values" 上的 topic
    t = pw.io.debezium.read(
        input_rdkafka_settings,
        topic_name="my_mongo_db.test_database.values",
        schema=InputSchema,
        autocommit_duration_ms=100,
    )

    # 计算求和(这部分与连接器无关,是纯 Pathway 操作)
    t = t.reduce(sum=pw.reducers.sum(t.value))

    # 结果同时写回 MongoDB 和文件系统
    pw.io.mongodb.write(
        t,
        connection_string="mongodb://mongodb:27017/?replicaSet=rs0",
        database="my_mongo_db",
        collection="sum_values",
    )
    pw.io.csv.write(t, "output_stream.csv")

    # 启动计算
    pw.run()

这段代码对应仓库中的完整示例 sum.py

整体架构与 Docker 五件套

整个链路涉及五部分:MongoDB、Debezium、Kafka、ZooKeeper、Pathway。从零手工配置这五个组件门槛较高,教程使用 Docker Compose 将其打包为一键部署。一个 Compose 文件的组织形式大致如下,每个应用称为一个 service,buildenvironmentvolumes 等参数分别定义容器构建方式、环境变量与挂载卷:

version: "3.7"
services:
  mongodb:
    build:
    environment:
    volumes:
  kafka:
    build:
    ...

注意 mongodb 在这里是 service 名,它实际使用什么镜像由 build/image 参数决定。数据流方向为:values 集合的写入 → Debezium 捕获 oplog 变更 → Kafka topic → Pathway 消费并求和 → 回写 sum_values 集合。其中 sum_values 无需手动创建,输出连接器会在首次写入时自动建库建集合。

完整编排文件见 docker-compose.yml,下面按组件逐一展开。

MongoDB:必须开启复制集

  mongodb:
    container_name: mongodb
    image: mongo
    command: ["--replSet", "rs0", "--port", "27017"]
    healthcheck:
      test: echo "try { rs.status() } catch (err) { rs.initiate({_id:'rs0',members:[{_id:0,host:'mongodb:27017'}]}) }" | mongosh --port 27017 --quiet
      interval: 5s
      timeout: 120s

配置要点:

  • container_name:容器名;
  • image:官方 MongoDB Docker 镜像;
  • command--replSet rs0 启用复制集并命名 rs0--port 27017 指定端口。Debezium 的 MongoDB 连接器依赖 oplog,而 oplog 只有复制集才提供,所以这一步不可省略(这与源码文档中 pw.io.mongodb.read 对复制集的强制要求一致,见 mongodb 连接器);
  • healthcheck:MongoDB 启动后复制集尚未初始化,需执行 rs.initiate。该 healthcheck 每 5 秒重试一次 rs.status(),失败则自动以单节点初始化 rs0,直到数据库就绪、复制集建立完成。

输入数据流器:向 MongoDB 持续灌数据

数据库跑起来后,用一个 Python 脚本作为输入流,循环从 1 到 1000、每隔 500 毫秒向输入集合插入一条数据:

import sys
import time

from pymongo import MongoClient

if __name__ == "__main__":
    client = MongoClient("mongodb://mongodb:27017/?replicaSet=rs0")
    db = client["test_database"]
    collection = db["values"]
    for i in range(1, 1001):
        collection.insert_one({"value": i})
        print("Insert:", i, file=sys.stderr)
        time.sleep(0.5)

对应 Compose 片段只有关键的依赖声明:

  streamer:
    container_name: mongodb-streamer
    build:
      context: .
      dockerfile: ./data-streaming/Dockerfile
    depends_on: [mongodb]

depends_on: [mongodb] 保证流器等待 MongoDB 完全初始化后再开始插入,避免向未就绪的数据库写入。该脚本对应 streamer.py

ZooKeeper 与 Kafka

Debezium Connect 需要 Kafka 作为消息总线,而该版本的 Confluent 镜像使用 ZooKeeper 做协调。由于 Docker Compose 中所有容器共享网络,可以直接用 service 名(如 zookeeper:2181)互相寻址:

  zookeeper:
    container_name: db_tuto_zookeeper
    image: confluentinc/cp-zookeeper:5.5.3
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
  kafka:
    container_name: db_tuto_kafka
    image: confluentinc/cp-enterprise-kafka:5.5.3
    depends_on: [zookeeper]
    environment:
      KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181"
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
      KAFKA_BROKER_ID: 1
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_JMX_PORT: 9991

其中 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 决定了其他容器(包括 Pathway 的 bootstrap.servers)应以 kafka:9092 访问 Kafka。

Debezium:通过 HTTP 注册 MongoDB 连接器

  debezium:
    container_name: db_tuto_debezium
    image: debezium/connect:1.4
    environment:
      BOOTSTRAP_SERVERS: kafka:9092
      GROUP_ID: 1
      CONFIG_STORAGE_TOPIC: connect_configs
      OFFSET_STORAGE_TOPIC: connect_offsets
    volumes:
      - ./debezium/connector.sh:/kafka/connector.sh
    depends_on: [kafka]
    ports:
      - 8083:8083

Debezium 容器的核心配置是一个挂载进去的注册脚本 connector.sh,它通过 REST 接口向 debezium:8083 发 HTTP 请求创建连接器,并轮询重试直到返回 201:

#!/bin/bash

while true; do
  http_code=$(curl -o /dev/null -w "%{http_code}" -H 'Content-Type: application/json' debezium:8083/connectors --data '{
    "name": "values-connector",
    "config": {
      "connector.class": "io.debezium.connector.mongodb.MongoDbConnector",
      "mongodb.hosts": "rs0/mongodb:27017",
      "mongodb.name": "my_mongo_db",
      "database.include.list": "test_database",
      "database.history.kafka.bootstrap.servers": "kafka:9092",
      "database.history.kafka.topic": "dbhistory.mongo"
    }
  }')
  if [ "$http_code" -eq 201 ]; then
    echo "Debezium connector has been created successfully"
    break
  else
    echo "Retrying Debezium connection creation in 1 second..."
    sleep 1
  fi
done

请求体中的关键参数:connector.class 指定 MongoDB 连接器类型;mongodb.hosts 使用 rs0/mongodb:27017 的形式指向复制集;mongodb.name 指定逻辑库名(它参与 topic 命名,这就是后面 Pathway 侧 topic 叫 my_mongo_db.test_database.values 的原因);database.include.list 限定只捕获 test_databasedatabase.history.kafka.* 则把数据库 schema 历史事件也落到 Kafka。完整的 MongoDB 连接器参数以 Debezium 官方文档为准。

⚠️ 该脚本需要可执行权限(Makefile 的 build 目标会先执行 chmod +x)。

Pathway 容器

  pathway:
    container_name: db_tuto_pathway
    build:
      context: .
      dockerfile: ./pathway-src/Dockerfile
    depends_on: [kafka, mongodb]

Pathway 容器的 Dockerfile 很简单,基于 Python 镜像 pip 安装 Pathway 并把入口脚本 sum.py 拷入容器:

FROM --platform=linux/x86_64 python:3.10

RUN pip install -U pathway
COPY ./pathway-src/sum.py sum.py

CMD ["python", "-u", "sum.py"]

⚠️ 出于兼容性考虑这里固定使用 x86_64 平台。sum.py 是容器入口点,容器启动即执行,脚本退出容器即停止。

Makefile:三个命令驱动整个实验

Makefile 把编排动作收敛为三个目标:

SERVICE_NAME_PATHWAY = pathway


build:
	chmod +x ./debezium/connector.sh
	docker compose up -d
	docker compose exec debezium ./connector.sh

stop:
	docker compose down -v

connect-pathway:
	docker compose exec $(SERVICE_NAME_PATHWAY) bash
  • make(默认执行 build):赋予注册脚本执行权限 → 后台启动全部容器 → 进入 debezium 容器执行 connector.sh 完成连接器注册;
  • make stop:拆除并清理容器与卷;
  • make connect-pathway:进入 Pathway 容器的 bash,可以查看文件系统连接器生成的 output_stream.csv

Debezium 输入连接器:pw.io.debezium.read

输入侧要求 topic 中的消息必须是 Debezium 封装格式;每条更新是原子的,会触发 Pathway 计算图中对应的一次更新。Debezium 连接器只工作在流式(streaming)模式,并且一次只监听一个 topic

pw.io.debezium.read 的完整参数如下(结合源码 docstring 补全了教程未提及的两项):

参数 说明
rdkafka_settings 连接接收 Debezium 消息的 Kafka 实例的设置,遵循 librdkafka 配置格式(如 bootstrap.serversgroup.idauto.offset.reset 等)
topic_name 监听的 topic 名称
db_type 事件来源的数据库类型,枚举 DebeziumDBType,可选 POSTGRES/MONGO_DB;源码中默认值为 POSTGRES(见 engine.pyi),因此从源码结构看,读取 MongoDB 主题时应显式传入 MONGO_DB。仓库中的示例未显式设置该参数,可视为早期用法,建议在新代码中显式声明
schema 定义 Pathway 表结构的 pw.Schema,指定列名、类型与主键
autocommit_duration_ms 两次提交之间的最大间隔(毫秒)。每经过该时长,连接器收到的更新会被提交并推入 Pathway 计算图;源码默认值为 1500
name 连接器的唯一名称,用于日志与监控面板;启用持久化时作为存储连接器进度快照的名称
max_backlog_size 限制任一时刻被读取并保留在处理中的条目数,达到上限即暂停读取。适合处理“初始突发大量数据”的大数据源,避免内存尖峰

使用示例(即 sum.py 中的读端):

class InputSchema(pw.Schema):
    value: int

t = pw.io.debezium.read(
    input_rdkafka_settings,
    topic_name="my_mongo_db.test_database.values",
    schema=InputSchema,
    autocommit_duration_ms=100,
)

从源码实现看,read 内部构造了一个 storage_type="kafka" 的 DataStorage 与 format_type="debezium" 的 DataFormat,再交给 table_from_datasource 生成 Pathway 表(见 debezium 连接器实现)——这印证了前文的说法:Debezium 连接器底层就是 Kafka 存储 + Debezium 消息格式解析。按 Debezium 的机制,消费开始时它会先发送集合的全量快照,随后持续流式推送 insert/update/delete 变更。

MongoDB 输出连接器:pw.io.mongodb.write

输出连接器把表 t 的更新追加到指定 MongoDB 集合。教程列出的基础参数为:table(要写入的表)、connection_string(MongoDB 连接串)、database(数据库名)、collection(集合名)、max_batch_size(单批最大插入条数)。结合 mongodb 连接器源码 还可以看到两个进阶参数:

参数 说明
output_table_type 输出形态,默认 "stream_of_changes"(保留完整变更历史);可选 "snapshot"(只维护表的当前状态,以 _id 定位行并直接更新)
sort_by 指定列引用列表,使每个 minibatch 内按这些列的元组字典序升序输出
name 连接器唯一名称,用于日志与监控
pw.io.mongodb.write(
    t,
    connection_string="mongodb://mongodb:27017/",
    database="my_mongo_db",
    collection="sum_values",
)

默认 stream_of_changes 模式下,表 t 每被更新一次,变更就会追加到 sum_values 集合,输出集合包含 t 的全部列,外加两个系统字段:time 标识该变更所在的事务性 mini-batchdiff 描述变更性质——diff = 1 表示插入,diff = -1 表示删除;一次“更新”在同一 mini-batch 中表现为两条事件:先以 diff = -1 移除旧值,再以 diff = 1 插入新值。集合(乃至数据库)不存在时会在首次写入时自动创建。

两个由源码确认的约束需要注意:输入表不能包含名为 _id 的列(MongoDB 的文档主键保留名);在 stream_of_changes 模式下也不能使用 timediff 作为列名,违反时会在构建阶段抛出 ValueError(见 write 实现的列名校验)。本例表只有 valuesum 两列,不触发这些约束。

完整示例:项目结构与运行验证

整个示例项目的目录结构如下:

.
├── data-streaming/
│   ├── streamer.py
│   └── Dockerfile
├── debezium/
│   └── connector.sh
├── pathway-src/
│   ├── Dockerfile
│   └── sum.py
├── docker-compose.yml
└── Makefile

以上各文件在仓库中均位于 debezium-mongodb-example 目录下,与本文引用的路径一一对应。

运行步骤

  1. 在项目根目录执行 make——启动全部容器、初始化复制集、注册 Debezium 连接器,随后 streamer 开始向 test_database.values 插入数据;
  2. 每次向 values 集合写入新条目,都会触发 Debezium 捕获 → Kafka topic → Pathway 求和 → 回写 sum_values 并把结果同步到 output_stream.csv
  3. make connect-pathway 进入 Pathway 容器查看 CSV 文件;
  4. 进入 MongoDB 容器实时观察回写结果:
docker-compose exec mongodb mongosh

连接后切换到目标库,查看最新的求和记录:

use my_mongo_db
db["sum_values"].find().sort({ time: -1 }).pretty()

该值由 Pathway 实时持续更新——time 单调递增的文档序列直观展示了 Stream of Updates 机制下每次 mini-batch 的产出。

小结

这条教程展示了一套可复制的 MongoDB 实时处理范式:复制集是 CDC 的前提,Debezium 负责把 oplog 变更转化为 Kafka 上的结构化消息,Pathway 侧仅用 pw.io.debezium.read 一个连接器即可获得“全量快照 + 增量流”的完整输入表,而 pw.io.mongodb.write 则把结果以带 time/diff 元数据的变更流形式落回数据库。整个链路五容器全部由一个 docker-compose.yml 与一个 Makefile 驱动,sum.py 中真正的计算逻辑只有一行 reduce。将 db_typemax_backlog_sizeoutput_table_type 等源码级参数纳入考虑后,这套架构可以平滑迁移到已有的 MongoDB/Debezium 基础设施上,仅保留输入与输出两段连接器代码即可复用。

登录后查看全文
热门项目推荐
相关项目推荐