Airbyte source-postgres 本地 CDC 端到端测试指南:用 source-postgres-e2e-cdc-tests Skill 复现逻辑解码复制行为
2026/9/23 12:17:46 网站建设 项目流程

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=logicalmax_replication_slots=10max_wal_senders=10启动,镜像固定为postgres:16(而非latest,以保证大版本行为稳定,可用BACKEND_IMAGE覆盖)。这些参数正是逻辑复制生效的前提:

参数作用
wal_levellogical启用逻辑解码所需的信息写入 WAL
max_replication_slots10允许的最大复制槽数量,供airbyte_slot使用
max_wal_senders10允许的最大 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 只订阅usersCREATE 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、uvjq,以及一份 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.txtstderr.txtreport.mdreport.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 的核心:

字段说明
methodCDC声明复制方式为逻辑解码 CDC
replication_slotairbyte_slot与 SQL fixture 中创建的复制槽对应
publicationairbyte_publication与 fixture 中创建的发布对应
initial_waiting_seconds120初始等待时间(秒),用于等待复制槽就绪

为什么必须显式传--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_idminimum_generation_idsync_id)和 destination 元数据(destination_object_nameinclude_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 为stdoutstderrany;命令名为speccheckdiscoverread。任一期望失败都会以非零退出码结束并进入生成的 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

  1. 基线 read:应用 fixture,跑一次完整 read(断言 ≥3 条记录、≥1 条 STATE);
  2. 提取状态:用共享的 extract-state.py 从$REPRO_OUT/<step>/read/stdout.txt中提取最新一条协议级 STATE 消息(该工具与引擎无关);
  3. 变更数据:通过通用apply-sql.sh应用 mutate-insert-users.sql,即插入一行('dave', 'dave@example.com')
  4. 带状态回放:以--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 复现用例

官方给出的创作流程如下:

  1. 添加幂等 SQL fixturefixtures/sql/下。CDC 设置必须放在被配置的数据库内,且不要添加按运行的服务端配置(逻辑复制参数在通用后端启动时已确立,case 中除了 SQL fixture 里的 publication 与复制槽,无需额外 CDC 设置);
  2. 需要新的流形状时,在fixtures/catalogs/下添加 CDC 感知目录。使用固定版本连接器discover输出的源定义游标;
  3. 添加cases/<name>.sh:以set -euo pipefail开头,带${VERSION:-3.8.5}默认值,调用通用 Skill 的scripts/run.sh,传入--config-template--catalog--fixture--keep-backend及相关的--expect-*flags;
  4. 多阶段用例仿照cases/incremental-replay.sh的序列:基线read→ 用extract-state.py$REPRO_OUT/<step>/read/stdout.txt提取 STATE → 用通用apply-sql.sh变更数据 → 以--skip-fixtures --state=PATH再跑;
  5. 验证:对已发布镜像或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.modeprefer;非 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_passwordBACKEND_PASSWORD覆盖),初始库默认test_dbBACKEND_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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询