Pathway 数据库连接器实战:用 Debezium CDC 与 PostgreSQL 构建实时数据管道
本文基于 Pathway 官方教程 Database connectors 展开,完整讲解如何用 pw.io.debezium.read 和 pw.io.postgres.write 两个连接器,把 PostgreSQL 的变更(CDC)实时摄入 Pathway 流式计算图并写回数据库。读完后你将掌握:五容器(PostgreSQL、Debezium、Kafka、ZooKeeper、Pathway)的 Docker Compose 部署方案、Debezium 输入连接器与 PostgreSQL 输出连接器的全部参数、以及一份可直接运行的完整示例(仓库内 debezium-postgres-example)。
1. 场景与核心思路
传统数据库(如 PostgreSQL)并非为流式场景设计:要监控数据库并把变更流式化,需要变更数据捕获(CDC)机制。该方案的组合是:
- 输入:Debezium 捕获 PostgreSQL 表
values的变更,经 Kafka 以 Debezium 消息格式推送;Pathway 用pw.io.debezium.read消费该 topic; - 处理:Pathway 计算增量总和;
- 输出:用
pw.io.postgres.write把结果流写回 PostgreSQL 表sum_table。
补充:如果你的数据库是 MongoDB,Pathway 提供同构的 MongoDB CDC 教程(结构相似、配置步骤不同);如果使用的是 NeonDB 这类 serverless、PostgreSQL 兼容的数据库,则可以不走 Debezium,直接用 Pathway 的原生 PostgreSQL 连接器
pw.io.postgres.read(streaming 模式基于逻辑复制槽)完成 CDC 读与写。原生读取连接器的实现与许可限制见下文第 5 节。
1.1 最短可用代码
以一张只有 value 一列的表 values 为例,对每个新值实时求和并写回 sum_table:
import pathway as pw
# Debezium settings
input_rdkafka_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
"group.id": "0",
"session.timeout.ms": "6000",
}
# PostgreSQL settings
output_postgres_settings = {
"host": "postgres",
"port": "5432",
"dbname": "values_db",
"user": "user",
"password": "password",
}
# 定义表结构:列名 + 类型
class InputSchema(pw.Schema):
value: int
# 用 Debezium 连接器监听 "postgres.public.values" topic
t = pw.io.debezium.read(
input_rdkafka_settings,
topic_name="postgres.public.values",
schema=InputSchema,
)
# 计算总和(这部分与连接器无关)
t = t.reduce(sum=pw.reducers.sum(t.value))
# 用 PostgreSQL 连接器把结果流写回数据库
pw.io.postgres.write(t, output_postgres_settings, "sum_table")
# 启动计算
pw.run()
完整可运行版本见 pathway-src/sum.py。
2. 架构:五个容器如何协作
整条链路包含五块拼图:
- PostgreSQL:两张表——被持续插入的输入表
values,以及被 Pathway 周期性更新的输出表sum_table; - Debezium:CDC 实例,捕获
values表的变更; - Kafka(依赖 ZooKeeper):Debezium 更新先发到 Kafka,Kafka 把消息传播给 Pathway;
- Pathway:从 Kafka 接收更新、计算、把和写回 PostgreSQL。
从零安装这些组件很繁琐,因此示例使用 Docker Compose 一键拉起全部容器。容器相当于一个独立应用的轻量虚拟环境(例如 PostgreSQL 容器内就是一个最小化的、可直接使用的 PostgreSQL 发行版)。主容器名即服务名,容器间共享同一网络,相互访问直接以服务名寻址(连 ZooKeeper 只写 "zookeeper:2181")。
2.1 PostgreSQL 服务
仓库中 docker-compose.yml 的 postgres 部分:
postgres:
container_name: db_tuto_postgres
image: debezium/postgres:13
environment:
- POSTGRES_USER=user
- POSTGRES_PASSWORD=password
- POSTGRES_DB=values_db
- PGPASSWORD=password
volumes:
- ./sql/init-db.sql:/docker-entrypoint-initdb.d/init-db.sql
- ./sql/update_db.sh:/update_db.sh
要点:
- 容器首次启动时会执行
/docker-entrypoint-initdb.d/下的脚本初始化数据库。这里通过 volume 把 init-db.sql 映射进去创建两张表; - update_db.sh 被映射到容器根目录(而非 initdb 目录),因为需要手动执行来持续插入数据、制造输入流。该脚本需有可执行权限。
仓库实际的建表脚本:
CREATE TABLE IF NOT EXISTS values (
value integer NOT NULL
);
CREATE TABLE IF NOT EXISTS sum_table (
id SERIAL PRIMARY KEY,
sum BIGINT NOT NULL,
time BIGINT NOT NULL,
diff SMALLINT NOT NULL
);
注意 sum_table 除了业务列 sum,还必须包含 time(BIGINT)与 diff(SMALLINT)两个元数据列——原因见第 4 节。仓库版本还额外加了 id SERIAL PRIMARY KEY,便于按写入顺序追溯每次更新。
数据生成脚本(循环插入 1000 条,间隔 0.5 秒):
#!/bin/bash
export PGPASSWORD='password'
sleep 3
for LOOP_ID in {1..1000}
do
psql -d values_db -U user -c "INSERT INTO values VALUES ($LOOP_ID);"
sleep 0.5
done
2.2 ZooKeeper 与 Kafka
Debezium 需要 Kafka,而这个版本的 Kafka 又依赖 ZooKeeper。示例选用专门定制的镜像,最大限度减少配置项:
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 通过 KAFKA_ZOOKEEPER_CONNECT 指向同网络的 ZooKeeper;KAFKA_ADVERTISED_LISTENERS 声明对外监听地址 PLAINTEXT://kafka:9092,这正是 Pathway 侧 bootstrap.servers 的取值。
2.3 Debezium 服务与连接器注册
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]
Debezium 容器本身只是 Kafka Connect 运行时,还需要向其 REST 端点(8083)POST 一次连接器定义,把 PostgreSQL 与 Kafka 接通。仓库中的 connector.sh 带了一个“轮询重试直到 HTTP 201”的循环,解决服务启动时序问题,比教程文档中的单次 curl 更健壮:
#!/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.postgresql.PostgresConnector",
"plugin.name": "pgoutput",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "user",
"database.password": "password",
"database.dbname" : "values_db",
"database.server.name": "postgres",
"table.include.list": "public.values",
"database.history.kafka.bootstrap.servers": "kafka:9092",
"database.history.kafka.topic": "schema-changes.inventory"
}
}')
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:PostgreSQL 源连接器;plugin.name: pgoutput:使用 PostgreSQL 自带的逻辑解码插件;table.include.list: public.values:只捕获values表——这决定了下游 topic 名将是postgres.public.values(即database.server.name.schema.table);database.history.kafka.*:schema 历史 topic,Debezium 用它保存被捕获表的 DDL 历史。
2.4 Pathway 服务
Pathway 不自带 Docker 镜像,用 Dockerfile 现场构建(见 pathway-src/Dockerfile):
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"]
pathway:
container_name: db_tuto_pathway
build:
context: .
dockerfile: ./pathway-src/Dockerfile
depends_on: [kafka, postgres]
sum.py 是容器入口,进程退出容器即停止。注意 --platform=linux/x86_64:由于兼容性考虑,该示例指定 x86_64 Linux 镜像(在 Apple Silicon 等 ARM 机器上会以 Rosetta 方式运行)。
2.5 Makefile:一键启动与停止
仓库实际的 Makefile(已改用 docker compose 子命令语法):
SERVICE_NAME_POSTGRES = postgres
SERVICE_NAME_PATHWAY = pathway
DB_NAME = values_db
USERNAME = user
build:
chmod +x ./debezium/connector.sh
chmod +x ./sql/update_db.sh
docker compose up -d
sleep 5
docker compose exec debezium ./connector.sh
docker compose exec postgres ./update_db.sh
stop:
docker compose down -v
# You may need to rename the image
docker rmi debezium-postgres-example-pathway:latest
connect-bash:
docker compose exec $(SERVICE_NAME_POSTGRES) bash
connect-psql:
docker compose exec $(SERVICE_NAME_POSTGRES) psql values_db -U user -W
connect-pathway:
docker compose exec $(SERVICE_NAME_PATHWAY) bash
make build(或make):拉起全部容器 → 等 5 秒 → 注册 Debezium 连接器 → 开始向values表持续插数;make stop:销毁容器并删除本地构建的 Pathway 镜像;make connect-psql:直接进入 postgres 容器查询结果。
3. Debezium 输入连接器 pw.io.debezium.read
数据流形态:输入必须是某个 topic 上的 Debezium 消息;每条收到的更新都是原子的,会触发 Pathway 计算图中对应表的更新。Debezium 连接器仅支持 streaming(流式)模式。
严格说,消息是“Debezium 格式”的:连接器本质上连接的是 Kafka,与普通 Kafka 连接器的唯一区别在于期望的消息格式不同。另注意:一个 Debezium 连接器只监听一个 topic——要读多张表就需要多个连接器实例。
结合 python/pathway/io/debezium/init.py 的函数签名,参数比教程文档更完整:
| 参数 | 说明 |
|---|---|
rdkafka_settings |
librdkafka 格式的连接配置(bootstrap.servers、group.id、session.timeout.ms、auto.offset.reset 等) |
topic_name |
要监听的 topic,如 postgres.public.values |
db_type |
事件来源数据库类型,默认 DebeziumDBType.POSTGRES |
schema |
pw.Schema 子类,定义列名、类型与主键 |
autocommit_duration_ms |
两次提交之间的最大间隔(毫秒),默认 1500;每隔该时长,连接器收到的更新被提交并推进计算图 |
name |
连接器唯一名,用于日志/监控面板;启用持久化时用作进度快照名 |
max_backlog_size |
限制同时保持在处理中的输入条目数,达到上限则暂停读取;适合初始数据突发很大的源,避免内存尖峰 |
debug_data |
调试模式下替代真实数据的静态数据 |
源码层面的实现细节(python/pathway/io/debezium/init.py):该连接器被组装为一个 GenericDataSource——底层 DataStorage 的 storage_type 是 "kafka",而 DataFormat 的 format_type 是 "debezium"。这印证了“Debezium 连接器 = Kafka 存储 + Debezium 消息格式解析”的判断;autocommit_duration_ms 则映射为数据源选项 commit_duration_ms。
典型用法:
class InputSchema(pw.Schema):
value: int
t = pw.io.debezium.read(
input_rdkafka_settings,
topic_name="postgres.public.values",
schema=InputSchema,
autocommit_duration_ms=100
)
Debezium 先发送表的全量快照,再流式发送后续行级变更(insert / update / delete),因此连接器启动时 t 会包含整表快照,之后持续追增量。
4. PostgreSQL 输出连接器 pw.io.postgres.write
输出表结构要求(stream of changes 模式):目标表必须包含 Pathway 表 t 的全部列,外加两个元数据列:
time:该次更新所在事务 minibatch 的时间戳(毫秒级,建议BIGINT,因为常规取值会超出 32 位范围);diff:1表示插入,-1表示删除。一次“更新”在流式语义中即旧值删除 + 新值插入同时发生。
因此本例的 sum_table 需要:
CREATE TABLE IF NOT EXISTS sum_table (
sum BIGINT NOT NULL,
time BIGINT NOT NULL,
diff SMALLINT NOT NULL
);
⚠️ 默认情况下目标表需预先创建,Pathway 不负责建表(可用 init_mode 参数改变这一点,见下)。
pw.io.postgres.write(t, output_postgres_settings, "sum_table")
结合 python/pathway/io/postgres/init.py 的签名,可用参数有:
| 参数 | 说明 |
|---|---|
table / table_name |
要写入的 Pathway 表;目标 PostgreSQL 表名(会做标识符引用,连字符、大小写、保留字可原样通过) |
postgres_settings |
libpq 风格的 key-value 连接参数字典,host/port/dbname/user/password 等;所有键值对被空格拼接成 libpq 连接串 |
schema_name |
目标表所在 schema,默认 "public" |
output_table_type |
"stream_of_changes"(默认,追加 time/diff 两列)或 "snapshot"(只保留当前状态、原子刷新,无额外列) |
primary_key |
snapshot 模式必填:目标表主键列的 ColumnReference 列表,底层据此生成 INSERT ... ON CONFLICT (...) DO UPDATE |
init_mode |
"default" / "create_if_not_exists"(自动建表)/ "replace"(重建表) |
max_batch_size |
单个事务内允许提交的最大条目数 |
name |
连接器唯一名,同时作为连接串中 application_name(pathway:<name>)的一部分,便于在 pg_stat_activity 中过滤 |
sort_by |
每个 minibatch 内按给定列升序输出 |
值得说明的两个实现细节(均可在 python/pathway/io/postgres/init.py 源码中确认):
- 连接串与超时默认值:
_augment_postgres_settings会为未显式给出的键注入保守的 TCP keepalive 默认值(keepalives_idle=300、keepalives_interval=30、keepalives_count=3、tcp_user_timeout=300000、connect_timeout=30等),使一个失联的 Pathway 进程在几分钟内被 PostgreSQL 检出,而不是等待操作系统默认的约两小时;用户显式传入的同名键总是优先。 - 参数校验很严格:stream 模式下若你的 Pathway schema 里恰好有名为
time/diff的业务列会直接报错(避免与元数据列冲突);snapshot 模式未给主键、主键列可空、主键引用了表中不存在的列等,都会在调用时抛出带修复建议的ValueError,而不是等到建表后才失败。
5. 延伸阅读:原生 PostgreSQL 读取连接器 pw.io.postgres.read
除了 Debezium 路线,Pathway 还提供原生 CDC 读取连接器(实现见 python/pathway/io/postgres/init.py,docstring 注明该模块仅在 Scale/Enterprise 许可下可用)。它支持两种模式:
mode="static":一次性读表后停止;mode="streaming":使用pgoutput逻辑解码插件创建临时复制槽(export snapshot),先快照、再消费 WAL 增量;槽随连接关闭自动删除,无需手工管理。streaming 模式要求服务端wal_level = logical并预先创建 publication:
CREATE PUBLICATION users_pub FOR TABLE users;
table = pw.io.postgres.read(
postgres_settings,
table_name="users",
schema=UsersSchema,
mode="streaming",
publication_name="users_pub",
)
对只需 CDC 读、不想再部署 Kafka + Debezium + ZooKeeper 三件套的场景,这是一条明显更轻的路线(NeonDB 集成即基于此连接器);而 Debezium 路线的优势在于与既有 Debezium/Kafka 基础设施复用、且支持更多源数据库。
6. 完整示例运行与结果验证
示例项目结构(与 examples/projects/debezium-postgres-example 一致):
examples/projects/debezium-postgres-example/
├── debezium/
│ └── connector.sh
├── pathway-src/
│ ├── Dockerfile
│ └── sum.py
├── sql/
│ ├── init-db.sql
│ └── update_db.sh
├── docker-compose.yml
└── Makefile
其中 sum.py 是容器入口,内容(仓库实际版本,还附带一段 t.debug 调试输出与 CSV 落盘,可忽略):
import pathway as pw
input_rdkafka_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
"group.id": "0",
"session.timeout.ms": "6000",
"auto.offset.reset": "earliest",
}
output_postgres_settings = {
"host": "postgres",
"port": "5432",
"dbname": "values_db",
"user": "user",
"password": "password",
}
class InputSchema(pw.Schema):
value: int
t = pw.io.debezium.read(
input_rdkafka_settings,
topic_name="postgres.public.values",
schema=InputSchema,
autocommit_duration_ms=100,
)
t = t.reduce(sum=pw.reducers.sum(t.value))
pw.io.postgres.write(t, output_postgres_settings, "sum_table")
pw.run()
运行步骤:
- 在项目根目录执行
make build(等价于拉起五个容器、注册连接器、启动插数脚本)。pw.run()之后计算会持续运行直到进程被终止; values表每次插入都会经 Debezium 触发一次更新,Pathway 随即向sum_table写入一条增量记录;- 进入 PostgreSQL 容器查看实时结果:
docker compose exec postgres psql values_db -U user -W(或make connect-psql)。
查看当前总和(取最新一条变更):
SELECT sum FROM sum_table ORDER BY time DESC, diff DESC LIMIT 1;
该值由 Pathway 实时更新。查看最近 10 条变更记录:
SELECT * FROM sum_table ORDER BY time DESC, diff DESC LIMIT 10;
由于 diff=1/-1 分别代表插入/删除,按 time、diff 降序取首行即可还原“当前累计状态”;这也是 stream of changes 输出表的通用解读方式。
7. 要点回顾
- 输入:
pw.io.debezium.read(rdkafka_settings, topic_name, schema=..., autocommit_duration_ms=...),仅流式模式、单 topic; - 输出:
pw.io.postgres.write(t, postgres_settings, table_name),stream of changes 模式要求目标表含time(BIGINT)与diff(SMALLINT/INTEGER)列;需要“只留当前态”时切换output_table_type="snapshot"并指定primary_key; - 部署:五个容器由一个
docker-compose.yml定义,Makefile封装启动/停止/连接三步; - 更轻的替代:若无需复用 Kafka 基础设施,可评估原生
pw.io.postgres.read(streaming 模式基于临时复制槽,需要 Scale/Enterprise 许可与wal_level=logical及 publication)。
以上所有可运行文件均可在仓库 examples/projects/debezium-postgres-example 中逐一对应查看。
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 StartedRust0624
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