Airbyte source-postgres 本地 CDC 端到端测试指南:用 source-postgres-e2e-cdc-tests Skill 复现逻辑解码复制行为
【免费下载链接】airbyteOpen-source data movement for ELT pipelines and AI agents — from APIs, databases & files to warehouses, lakes, and AI applications. Both self-hosted and Cloud.项目地址: https://gitcode.com/gh_mirrors/ai/airbyte
本指南以 source-postgres-e2e-cdc-tests SKILL.md 为核心,系统讲解如何在本地搭建 PostgreSQL 逻辑解码(logical decoding)CDC 测试环境,运行初始加载(initial load)与增量回放(incremental replay)两类冒烟用例,并手把手带你编写幂等的 SQL fixture 与基于声明式断言的 case 脚本,用于复现source-postgres的 CDC 模式缺陷、验证修复后的连接器镜像。读完你将掌握一套可在 Docker 本地一键运行的 CDC 回归验证流程。
1. 这套 Skill 解决什么问题
source-postgres-e2e-cdc-tests是 Airbyte 仓库中为source-postgres连接器提供的本地 CDC 测试工具集(Skill)。它的定位非常聚焦:在你本机的 Docker 环境中,用一个开启了逻辑复制参数的 PostgreSQL 16 后端,配合 CDC 感知的配置、目录(catalog)和 SQL fixture,完整跑通 PostgreSQL 逻辑解码(logical decoding)CDC 行为,从而:
- 本地复现
source-postgres在 CDC 模式下的 bug; - 用已发布的连接器镜像或本地构建的
VERSION=dev镜像验证修复是否生效; - 运行 initial-load 与 incremental-replay 两个冒烟用例;
- 基于幂等 SQL fixture 与调用通用 Skill 的
scripts/run.sh的 case 脚本,编写全新的 PostgreSQL CDC 复现用例。
它的架构是分层组合的:引擎无关(engine-independent)的编排逻辑放在 airbyte-integrations/db-harness-lib/,通用 PostgreSQL Skill(source-postgres-e2e-tests)负责拉起 PostgreSQL 16 后端并运行连接器,而本 Skill 只补充 CDC 专属的配置、目录、SQL fixture 与 case 脚本。
source-postgres-e2e-cdc-tests/ ├── SKILL.md ├── cases/ │ ├── initial-load.sh # CDC 初始加载冒烟用例 │ └── incremental-replay.sh # read → mutate → read-with-state 冒烟用例 └── fixtures/ ├── configs/ │ └── cdc.template.json # CDC 复制配置 ├── catalogs/ │ └── users-cdc.json # 带 PostgreSQL CDC 元数据的 users 流 └── sql/ ├── 00-init-cdc.sql # publication、slot、表与数据行 └── mutate-insert-users.sql # 回放阶段之间插入的数据行2. 工作原理:逻辑解码、Publication 与 Replication Slot
PostgreSQL CDC 的核心机制是逻辑解码:通过一个publication(发布)和一个replication slot(复制槽),将 WAL 中的变更以逻辑变更流的形式输出。source-postgres连接器在 CDC 模式下正是订阅这个变更流来捕获数据。
关键点在于:通用 Skill 启动的 PostgreSQL 后端本身就带着本 Skill 所需的 WAL 与复制设置,因此本 Skill不会启动第二个后端,也不需要为每次运行单独配置服务端参数。参见 start-backend.sh 中容器启动参数:
docker run -d --rm \ --name "$BACKEND_NAME" \ -e POSTGRES_PASSWORD="$BACKEND_PASSWORD" \ -e POSTGRES_DB="$BACKEND_DB" \ -p "$BACKEND_PORT:5432" \ "$BACKEND_IMAGE" \ -c wal_level=logical \ -c max_replication_slots=10 \ -c max_wal_senders=10即后端以wal_level=logical、max_replication_slots=10、max_wal_senders=10启动,镜像固定为postgres:16(而非latest,以保证大版本行为稳定,可用BACKEND_IMAGE覆盖)。这些参数正是逻辑复制生效的前提:
| 参数 | 值 | 作用 |
|---|---|---|
wal_level | logical | 启用逻辑解码所需的信息写入 WAL |
max_replication_slots | 10 | 允许的最大复制槽数量,供airbyte_slot使用 |
max_wal_senders | 10 | 允许的最大 WAL 发送进程数 |
在数据库层面,CDC 的“零件”由 SQL fixture 建立。00-init-cdc.sql 展示了完整的初始化序列:
-- 幂等清理:若复制槽已存在则删除 DO $$ BEGIN IF EXISTS (SELECT 1 FROM pg_replication_slots WHERE slot_name = 'airbyte_slot') THEN PERFORM pg_drop_replication_slot('airbyte_slot'); END IF; END $$; DROP PUBLICATION IF EXISTS airbyte_publication; DROP TABLE IF EXISTS public.users; CREATE TABLE public.users ( id SERIAL PRIMARY KEY, name VARCHAR(64) NOT NULL, email VARCHAR(200) NOT NULL ); INSERT INTO public.users (name, email) VALUES ('alice', 'alice@example.com'), ('bob', 'bob@example.com'), ('carol', 'carol@example.com'); CREATE PUBLICATION airbyte_publication FOR TABLE public.users; SELECT pg_create_logical_replication_slot('airbyte_slot', 'pgoutput');注意两处细节:
- Publication 只订阅
users表:CREATE PUBLICATION airbyte_publication FOR TABLE public.users; - 复制槽采用
pgoutput插件:pg_create_logical_replication_slot('airbyte_slot', 'pgoutput'),这也是source-postgres逻辑解码默认使用的输出插件。
fixture 使用通用后端中的test_db数据库,原因很实际:PostgreSQL 无法 drop 掉apply-sql.sh当前正连接的数据库,所以只能在test_db内重置users表、publication 和逻辑复制槽。
3. 前置条件
与通用 Skill 一致:
- Docker、
uv、jq,以及一份 Airbyte 仓库的克隆; - 后端容器
source-postgres-db-backend,由 start-backend.sh 启动。
不需要GSM 或 Cloud 管理员凭证。这套工具只能用于本地测试,严禁对着客户连接或 Airbyte Cloud 运行。
uv的用途是安装airbyte-internal-ops工具(uv tool install airbyte-internal-ops可将airbyte-ops放入$PATH,或者每次调用前缀uvx airbyte-internal-ops),case 脚本最终通过airbyte-ops cloud connector regression-test对指定镜像执行协议命令。
4. 快速上手:运行两条冒烟用例
官方推荐的完整会话流程如下(可直接复制执行):
SKILL=airbyte-integrations/connectors/source-postgres/.agents/skills/source-postgres-e2e-cdc-tests GENERIC=airbyte-integrations/connectors/source-postgres/.agents/skills/source-postgres-e2e-tests LIB=airbyte-integrations/db-harness-lib export REPRO_OUT=/tmp/source-postgres-repro # 启动一次逻辑解码就绪的后端 "$GENERIC/scripts/start-backend.sh" # 运行任意一个冒烟用例(两者都传 --keep-backend) "$SKILL/cases/initial-load.sh" "$SKILL/cases/incremental-replay.sh" # 代码变更后使用本地构建的连接器镜像 VERSION=dev "$SKILL/cases/initial-load.sh" # 会话结束后的清理 BACKEND_NAME=source-postgres-db-backend \ "$LIB/scripts/stop-backend.sh"约定要点:
- 用例复用同一个通用后端并传递
--keep-backend:在会话开始处启动一次后端,运行所有用例,结束时再拆除,避免每个用例反复重建容器; - 默认
VERSION=3.8.5(当前该工具集使用的已发布镜像),本地构建镜像后可用VERSION=dev覆盖; - 后端拆除使用 db-harness-lib 的 stop-backend.sh,它是幂等的。
4.1 关于目标镜像的选择
对于已推送到 PR 分支的代码,优先使用已发布的目标镜像:通过 Airbyte Ops MCP 工具publish_connector_to_airbyte_registry发布预发布版本,然后把得到的<version>-preview.<7-char-sha>标签作为--test-version传入。具体流程见 db-harness-lib README 的 "Getting a target image" 章节。只有当代码不在已推送的 PR 分支上时,才使用--test-version=dev(或让run.sh通过:dockerBuildx本地构建)。
5. 执行引擎:从 case 脚本到 db-harness-lib
每个 case 都会调用通用 Skill 的引擎 shim(run.sh),再由它委托给 db-harness-lib 的 run.sh。这条调用链是理解整套工具的关键:
cases/<name>.sh └─> source-postgres-e2e-tests/scripts/run.sh(引擎 shim,导出环境变量) └─> db-harness-lib/scripts/run.sh(引擎无关编排) ├─ 应用 SQL fixture(apply-sql.sh) ├─ 渲染 config(render-config.sh,替换后端地址) ├─ 运行请求的连接器命令(run-protocol-cmd.sh 调用 airbyte-ops) ├─ 将产物存储到 $REPRO_OUT/<step-name>/ └─ 强制执行 case 断言引擎 shim 的核心职责是导出引擎契约所需的环境变量(见 db-harness-lib/README.md 的 "Engine contract" 章节):
export CONNECTOR=source-postgres export ENGINE_SCRIPTS_DIR="$SKILL_DIR/scripts" export DEFAULT_CONFIG_TEMPLATE="$SKILL_DIR/fixtures/configs/base.template.json" export DEFAULT_FIXTURE="$SKILL_DIR/fixtures/sql/00-init-base.sql" export BACKEND_NAME="${BACKEND_NAME:-source-postgres-db-backend}"编排层负责:协议编排(spec → check → discover → read)、目录推导、配置渲染与状态提取。每次run.sh调用都会把产物写入$REPRO_OUT/<step-name>/下:渲染后的config.json、推导出的configured_catalog.json,以及每个协议命令一个子目录(spec/、check/、discover/、read/),里面存放该命令的stdout.txt、stderr.txt、report.md和report.html。
5.1 CDC 配置与目录为什么必须显式指定
cdc.template.json 是 CDC 模式的连接器配置模板:
{ "host": "source-postgres-db-backend", "port": 5432, "database": "test_db", "schemas": ["public"], "username": "postgres", "password": "test_password", "ssl_mode": { "mode": "prefer" }, "tunnel_method": { "tunnel_method": "NO_TUNNEL" }, "replication_method": { "method": "CDC", "replication_slot": "airbyte_slot", "publication": "airbyte_publication", "initial_waiting_seconds": 120 } }replication_method段是 CDC 的核心:
| 字段 | 值 | 说明 |
|---|---|---|
method | CDC | 声明复制方式为逻辑解码 CDC |
replication_slot | airbyte_slot | 与 SQL fixture 中创建的复制槽对应 |
publication | airbyte_publication | 与 fixture 中创建的发布对应 |
initial_waiting_seconds | 120 | 初始等待时间(秒),用于等待复制槽就绪 |
为什么必须显式传--catalog?db-harness-lib/README.md 专门用一节解释了这一陷阱:当没有--catalog时,run.sh会从discover结果推导 read 目录,且--sync-mode默认是full_refresh、无游标。在 CDC 配置模板下,这样推导出的目录会配置出零个 CDC 流——连接器虽然仍会运行全局 CDC feed 并发出冷启动状态,但 read 根本不会执行 CDC,第二轮甚至可能以误导性的错误拒绝自己的状态(如Incumbent CDC state is invalid ... Saved offset no longer present)。在对比模式下,控制镜像与目标镜像以相同方式失败,看起来就像连接器既有 bug。因此:
只要配置里
replication_method.method == "CDC",就必须显式传--catalog=PATH(CDC Skill 自带,例如fixtures/catalogs/users-cdc.json),或者推导一个增量目录:--sync-mode=incremental --cursor-field=CURSOR --streams=TABLE1,TABLE2。若 CDC 配置遇到 full-refresh 推导目录,run.sh会拒绝执行 read 并以退出码 2 提示应传的 flags。
users-cdc.json 给出了一个带完整 CDC 元数据的目录示例,注意其中几个字段:
"source_defined_cursor": true, "default_cursor_field": ["_ab_cdc_lsn"], "source_defined_primary_key": [["id"]], "is_resumable": true, "is_file_based": false, "sync_mode": "incremental", "destination_sync_mode": "append_dedup", "cursor_field": ["_ab_cdc_lsn"], "generation_id": 0, "minimum_generation_id": 0, "sync_id": 0, "destination_object_name": "users", "include_files": false关于这些字段有一条重要的版本事实:PostgreSQL 3.8.5 会把_ab_cdc_lsn识别为源定义游标,且 discover 输出中不会包含较新的生成目录字段。但目录中仍然配置了is_file_based、generation 元数据(generation_id、minimum_generation_id、sync_id)和 destination 元数据(destination_object_name、include_files),因为固定版本镜像在初始化read时需要这些字段。
5.2 声明式断言:expect 系列 flags
与传统 driver 脚本手写grep -q '<substring>' || exit 1不同,run.sh接受声明式期望 flags,并在返回前强制执行:
| Flag | 作用 |
|---|---|
--expect-test=pass\|fail | 目标整体判定 |
--expect-control=pass\|fail | 对比模式下的控制镜像判定(需配合--control-version与--reset=fixture或--reset=backend) |
--min-records=N | 目标 read 必须至少包含 N 条RECORD消息 |
--min-states=N | 目标 read 必须至少包含 N 条STATE消息 |
--expect-match=[<command>:]<channel>:<regex>[:N] | 目标命令输出必须匹配该正则至少 N 次 |
--forbid-match=[<command>:]<channel>:<regex> | 目标命令输出不得匹配该正则 |
匹配断言使用指定命令的目标侧产物,未给命令前缀时默认read;接受的 channel 为stdout、stderr、any;命令名为spec、check、discover、read。任一期望失败都会以非零退出码结束并进入生成的 summary。
6. 两个冒烟用例的逐步拆解
6.1 初始 CDC 加载(initial-load)
initial-load.sh 完整内容:
#!/usr/bin/env bash set -euo pipefail HERE="$(cd "$(dirname "$0")" && pwd)" SKILL="$(cd "$HERE/.." && pwd)" GENERIC="$(cd "$SKILL/../source-postgres-e2e-tests" && pwd)" "$GENERIC/scripts/run.sh" \ --command=read \ --test-version="${VERSION:-3.8.5}" \ --step-name=cdc-initial-load \ --fixture="$SKILL/fixtures/sql/00-init-cdc.sql" \ --config-template="$SKILL/fixtures/configs/cdc.template.json" \ --catalog="$SKILL/fixtures/catalogs/users-cdc.json" \ --keep-backend \ --expect-test=pass \ --min-records=3 \ --min-states=1它做的事:应用00-init-cdc.sql(重建表、publication、复制槽并插入三行数据),再用cdc.template.json配置与users-cdc.json目录读取配置好的users流;断言 read 通过、至少 3 条记录、至少 1 条 STATE 消息。这验证了“逻辑解码就绪的后端 + CDC 配置 + 目录元数据 + 初始加载”四者能协同工作。
6.2 增量回放(incremental-replay)
incremental-replay.sh 展示了一个多阶段(multi-phase)用例的标准模式:
#!/usr/bin/env bash set -euo pipefail HERE="$(cd "$(dirname "$0")" && pwd)" SKILL="$(cd "$HERE/.." && pwd)" GENERIC="$(cd "$SKILL/../source-postgres-e2e-tests" && pwd)" REPO_ROOT="$(git -C "$HERE" rev-parse --show-toplevel)" LIB="$REPO_ROOT/airbyte-integrations/db-harness-lib" REPRO_OUT="${REPRO_OUT:-/tmp/source-postgres-repro}" VERSION="${VERSION:-3.8.5}" STEP_NAME="${STEP_NAME:-cdc-replay}" BASELINE_STATE="$REPRO_OUT/$STEP_NAME/state.json" export REPRO_OUT # 阶段一:基线 read "$GENERIC/scripts/run.sh" \ --command=read \ --test-version="$VERSION" \ --step-name="$STEP_NAME/baseline" \ --fixture="$SKILL/fixtures/sql/00-init-cdc.sql" \ --config-template="$SKILL/fixtures/configs/cdc.template.json" \ --catalog="$SKILL/fixtures/catalogs/users-cdc.json" \ --keep-backend \ --expect-test=pass \ --min-records=3 \ --min-states=1 # 提取 STATE 消息 mkdir -p "$(dirname "$BASELINE_STATE")" "$LIB/scripts/extract-state.py" \ "$REPRO_OUT/$STEP_NAME/baseline/read/stdout.txt" \ > "$BASELINE_STATE" # 插入新数据 "$GENERIC/scripts/apply-sql.sh" \ "$SKILL/fixtures/sql/mutate-insert-users.sql" # 阶段二:带状态回放 "$GENERIC/scripts/run.sh" \ --command=read \ --test-version="$VERSION" \ --step-name="$STEP_NAME/replay" \ --skip-fixtures \ --config-template="$SKILL/fixtures/configs/cdc.template.json" \ --catalog="$SKILL/fixtures/catalogs/users-cdc.json" \ --state="$BASELINE_STATE" \ --keep-backend \ --expect-test=pass \ --expect-match='stdout:dave@example\.com' \ --forbid-match='stdout:alice@example\.com' \ --min-records=1 \ --min-states=1流程可以概括为read → mutate → read-with-state:
- 基线 read:应用 fixture,跑一次完整 read(断言 ≥3 条记录、≥1 条 STATE);
- 提取状态:用共享的 extract-state.py 从
$REPRO_OUT/<step>/read/stdout.txt中提取最新一条协议级 STATE 消息(该工具与引擎无关); - 变更数据:通过通用
apply-sql.sh应用 mutate-insert-users.sql,即插入一行('dave', 'dave@example.com'); - 带状态回放:以
--skip-fixtures(不再重新应用 fixture)与--state=PATH再次 read,断言回放通过、输出包含dave@example.com且不包含alice@example.com。
这个用例证明了保存的 LSN 能在初始加载之后正确续读:回放只输出新增行而不重复输出旧行。--skip-fixtures在这里很关键——因为00-init-cdc.sql每次应用都会 drop 并重建表、publication 与airbyte_slot,会重置 fixture,所以第二阶段必须跳过。
注意:这两个是冒烟哨兵用例(smoke canaries),不是针对某个 bug 的复现。在本 Skill 出现之前,仓库中不存在任何 PostgreSQL CDC fixture 或按 bug 组织的复现示例。
7. 编写一个新的 CDC 复现用例
官方给出的创作流程如下:
- 添加幂等 SQL fixture到
fixtures/sql/下。CDC 设置必须放在被配置的数据库内,且不要添加按运行的服务端配置(逻辑复制参数在通用后端启动时已确立,case 中除了 SQL fixture 里的 publication 与复制槽,无需额外 CDC 设置); - 需要新的流形状时,在
fixtures/catalogs/下添加 CDC 感知目录。使用固定版本连接器discover输出的源定义游标; - 添加
cases/<name>.sh:以set -euo pipefail开头,带${VERSION:-3.8.5}默认值,调用通用 Skill 的scripts/run.sh,传入--config-template、--catalog、--fixture、--keep-backend及相关的--expect-*flags; - 多阶段用例仿照
cases/incremental-replay.sh的序列:基线read→ 用extract-state.py从$REPRO_OUT/<step>/read/stdout.txt提取 STATE → 用通用apply-sql.sh变更数据 → 以--skip-fixtures --state=PATH再跑; - 验证:对已发布镜像或
VERSION=dev运行 case,检查其产物,后端生命周期保持在通用 Skill 中管理。
还有一个实用的工程细节:通用apply-sql.sh通过 stdin 调用psql,所以 fixture 中可以包含多条 SQL 语句,幂等性设计完全可行。
8. 常见坑与注意事项
- 不要把本工具用于客户连接或 Airbyte Cloud——它只面向本地测试;
- CDC 配置必须配显式目录,否则
run.sh会拒绝执行 read(退出码 2)或产生零 CDC 流的误导性失败(详见上文第 5.1 节); - PostgreSQL 配置形状有讲究:本地明文后端下
ssl_mode.mode用prefer;非 CDC 标准同步的replication_method.method值是Standard; - 连接器通过 bridge IP 解析后端:本地 regression-test 命令目前不接受
--network,所以共享的 render-config.sh 会把后端地址替换进渲染后的配置(可用CONFIG_HOST_JQ覆盖 host 放置位置的 jq 表达式,例如.host = $h | .port = 5432); check失败也可能以零退出:连接器可能在退出码为 0 时仍发出失败的CONNECTION_STATUS,若复现依赖 check 时行为,要对状态消息做断言;- 后端容器名
source-postgres-db-backend硬编码在生命周期脚本中,仅在做并行测试隔离时用BACKEND_NAME=…覆盖,且不要使用客户连接名;密码默认test_password(BACKEND_PASSWORD覆盖),初始库默认test_db(BACKEND_DB覆盖); - 清理:
stop-backend.sh幂等,结束后可rm -rf "$REPRO_OUT"并unset REPRO_OUT。
9. 相关资源
- 本 Skill:source-postgres-e2e-cdc-tests/SKILL.md
- 通用 PostgreSQL Skill:source-postgres-e2e-tests/SKILL.md(后端生命周期、协议命令 sweep、
--control-version对比模式) - 引擎无关编排库:db-harness-lib/README.md(引擎契约、CDC 目录陷阱、脚本说明)
- 核心实现:00-init-cdc.sql、cdc.template.json、users-cdc.json
【免费下载链接】airbyteOpen-source data movement for ELT pipelines and AI agents — from APIs, databases & files to warehouses, lakes, and AI applications. Both self-hosted and Cloud.项目地址: https://gitcode.com/gh_mirrors/ai/airbyte
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考