尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Airbyte source-postgres 本地 CDC 端到端测试指南:用 source-postgres-e2e-cdc-tests Skill 复现逻辑解码复制行为

Airbyte source-postgres 本地 CDC 端到端测试指南:用 source-postgres-e2e-cdc-tests Skill 复现逻辑解码复制行为 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 decodingCDC 测试环境运行初始加载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 decodingCDC 行为从而本地复现source-postgres在 CDC 模式下的 bug用已发布的连接器镜像或本地构建的VERSIONdev镜像验证修复是否生效运行 initial-load 与 incremental-replay 两个冒烟用例基于幂等 SQL fixture 与调用通用 Skill 的scripts/run.sh的 case 脚本编写全新的 PostgreSQL CDC 复现用例。它的架构是分层组合的引擎无关engine-independent的编排逻辑放在 airbyte-integrations/db-harness-lib/通用 PostgreSQL Skillsource-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 SlotPostgreSQL 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_levellogical \ -c max_replication_slots10 \ -c max_wal_senders10即后端以wal_levellogical、max_replication_slots10、max_wal_senders10启动镜像固定为postgres:16而非latest以保证大版本行为稳定可用BACKEND_IMAGE覆盖。这些参数正是逻辑复制生效的前提参数值作用wal_levellogical启用逻辑解码所需的信息写入 WALmax_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, aliceexample.com), (bob, bobexample.com), (carol, carolexample.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-opscase 脚本最终通过airbyte-ops cloud connector regression-test对指定镜像执行协议命令。4. 快速上手运行两条冒烟用例官方推荐的完整会话流程如下可直接复制执行SKILLairbyte-integrations/connectors/source-postgres/.agents/skills/source-postgres-e2e-cdc-tests GENERICairbyte-integrations/connectors/source-postgres/.agents/skills/source-postgres-e2e-tests LIBairbyte-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 # 代码变更后使用本地构建的连接器镜像 VERSIONdev $SKILL/cases/initial-load.sh # 会话结束后的清理 BACKEND_NAMEsource-postgres-db-backend \ $LIB/scripts/stop-backend.sh约定要点用例复用同一个通用后端并传递--keep-backend在会话开始处启动一次后端运行所有用例结束时再拆除避免每个用例反复重建容器默认VERSION3.8.5当前该工具集使用的已发布镜像本地构建镜像后可用VERSIONdev覆盖后端拆除使用 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-versiondev或让run.sh通过:dockerBuildx本地构建。5. 执行引擎从 case 脚本到 db-harness-lib每个 case 都会调用通用 Skill 的引擎 shimrun.sh再由它委托给 db-harness-lib 的 run.sh。这条调用链是理解整套工具的关键cases/name.sh └─ source-postgres-e2e-tests/scripts/run.sh引擎 shim导出环境变量 └─ db-harness-lib/scripts/run.sh引擎无关编排 ├─ 应用 SQL fixtureapply-sql.sh ├─ 渲染 configrender-config.sh替换后端地址 ├─ 运行请求的连接器命令run-protocol-cmd.sh 调用 airbyte-ops ├─ 将产物存储到 $REPRO_OUT/step-name/ └─ 强制执行 case 断言引擎 shim 的核心职责是导出引擎契约所需的环境变量见 db-harness-lib/README.md 的 Engine contract 章节export CONNECTORsource-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 的核心字段值说明methodCDC声明复制方式为逻辑解码 CDCreplication_slotairbyte_slot与 SQL fixture 中创建的复制槽对应publicationairbyte_publication与 fixture 中创建的发布对应initial_waiting_seconds120初始等待时间秒用于等待复制槽就绪为什么必须显式传--catalogdb-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就必须显式传--catalogPATHCDC Skill 自带例如fixtures/catalogs/users-cdc.json或者推导一个增量目录--sync-modeincremental --cursor-fieldCURSOR --streamsTABLE1,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-testpass\|fail目标整体判定--expect-controlpass\|fail对比模式下的控制镜像判定需配合--control-version与--resetfixture或--resetbackend--min-recordsN目标 read 必须至少包含 N 条RECORD消息--min-statesN目标 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-loadinitial-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 \ --commandread \ --test-version${VERSION:-3.8.5} \ --step-namecdc-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-testpass \ --min-records3 \ --min-states1它做的事应用00-init-cdc.sql重建表、publication、复制槽并插入三行数据再用cdc.template.json配置与users-cdc.json目录读取配置好的users流断言 read 通过、至少 3 条记录、至少 1 条 STATE 消息。这验证了“逻辑解码就绪的后端 CDC 配置 目录元数据 初始加载”四者能协同工作。6.2 增量回放incremental-replayincremental-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 \ --commandread \ --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-testpass \ --min-records3 \ --min-states1 # 提取 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 \ --commandread \ --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-testpass \ --expect-matchstdout:daveexample\.com \ --forbid-matchstdout:aliceexample\.com \ --min-records1 \ --min-states1流程可以概括为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, daveexample.com)带状态回放以--skip-fixtures不再重新应用 fixture与--statePATH再次 read断言回放通过、输出包含daveexample.com且不包含aliceexample.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 --statePATH再跑验证对已发布镜像或VERSIONdev运行 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 5432check失败也可能以零退出连接器可能在退出码为 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. 相关资源本 Skillsource-postgres-e2e-cdc-tests/SKILL.md通用 PostgreSQL Skillsource-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),仅供参考
返回列表