Pathway × Debezium:用 Change Data Capture 实现 MongoDB 实时流处理与聚合回写
本文以 Pathway 仓库中的 MongoDB + Debezium 官方教程及其配套示例项目为主线,讲解如何为 MongoDB 这类非流式数据库搭建完整的变更数据捕获(CDC)链路:MongoDB 复制集 → Debezium → Kafka/ZooKeeper → Pathway,再将实时聚合结果通过 pw.io.mongodb.write 回写到 MongoDB。读完本文,你将能够独立部署这条五容器架构,理解 pw.io.debezium.read 与 pw.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,build、environment、volumes 等参数分别定义容器构建方式、环境变量与挂载卷:
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_database;database.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.servers、group.id、auto.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-batch,diff 描述变更性质——diff = 1 表示插入,diff = -1 表示删除;一次“更新”在同一 mini-batch 中表现为两条事件:先以 diff = -1 移除旧值,再以 diff = 1 插入新值。集合(乃至数据库)不存在时会在首次写入时自动创建。
两个由源码确认的约束需要注意:输入表不能包含名为 _id 的列(MongoDB 的文档主键保留名);在 stream_of_changes 模式下也不能使用 time、diff 作为列名,违反时会在构建阶段抛出 ValueError(见 write 实现的列名校验)。本例表只有 value、sum 两列,不触发这些约束。
完整示例:项目结构与运行验证
整个示例项目的目录结构如下:
.
├── data-streaming/
│ ├── streamer.py
│ └── Dockerfile
├── debezium/
│ └── connector.sh
├── pathway-src/
│ ├── Dockerfile
│ └── sum.py
├── docker-compose.yml
└── Makefile
以上各文件在仓库中均位于 debezium-mongodb-example 目录下,与本文引用的路径一一对应。
运行步骤:
- 在项目根目录执行
make——启动全部容器、初始化复制集、注册 Debezium 连接器,随后streamer开始向test_database.values插入数据; - 每次向
values集合写入新条目,都会触发 Debezium 捕获 → Kafka topic → Pathway 求和 → 回写sum_values并把结果同步到output_stream.csv; - 用
make connect-pathway进入 Pathway 容器查看 CSV 文件; - 进入 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_type、max_backlog_size、output_table_type 等源码级参数纳入考虑后,这套架构可以平滑迁移到已有的 MongoDB/Debezium 基础设施上,仅保留输入与输出两段连接器代码即可复用。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0627
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00