首页
/ Pathway 数据库连接器实战:用 Debezium CDC 与 PostgreSQL 构建实时数据管道

Pathway 数据库连接器实战:用 Debezium CDC 与 PostgreSQL 构建实时数据管道

2026-09-06 12:30:37作者:凌朦慧Richard

本文基于 Pathway 官方教程 Database connectors 展开,完整讲解如何用 pw.io.debezium.readpw.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. 架构:五个容器如何协作

整条链路包含五块拼图:

  1. PostgreSQL:两张表——被持续插入的输入表 values,以及被 Pathway 周期性更新的输出表 sum_table
  2. Debezium:CDC 实例,捕获 values 表的变更;
  3. Kafka(依赖 ZooKeeper):Debezium 更新先发到 Kafka,Kafka 把消息传播给 Pathway;
  4. 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.serversgroup.idsession.timeout.msauto.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——底层 DataStoragestorage_type"kafka",而 DataFormatformat_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 位范围);
  • diff1 表示插入,-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_namepathway:<name>)的一部分,便于在 pg_stat_activity 中过滤
sort_by 每个 minibatch 内按给定列升序输出

值得说明的两个实现细节(均可在 python/pathway/io/postgres/init.py 源码中确认):

  • 连接串与超时默认值_augment_postgres_settings 会为未显式给出的键注入保守的 TCP keepalive 默认值(keepalives_idle=300keepalives_interval=30keepalives_count=3tcp_user_timeout=300000connect_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()

运行步骤

  1. 在项目根目录执行 make build(等价于拉起五个容器、注册连接器、启动插数脚本)。pw.run() 之后计算会持续运行直到进程被终止;
  2. values 表每次插入都会经 Debezium 触发一次更新,Pathway 随即向 sum_table 写入一条增量记录;
  3. 进入 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 分别代表插入/删除,按 timediff 降序取首行即可还原“当前累计状态”;这也是 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 中逐一对应查看。

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