资讯动态

SeaTunnel PostgreSQL CDC 源连接器完全指南:从快照到 WAL 的增量同步实战

发布时间:2026/9/20 3:10:31 来源:尧图企业网站定制
数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载PostgreSQL CDCChange Data Capture源连接器是 SeaTunnel 中基于逻辑复制槽技术实现 PostgreSQL 全量 增量数据同步的核心组件。本文以官方文档为骨架结合seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres模块源码系统讲解环境准备、连接器配置、快照与 WAL 读取机制、Schema 演进、无主键表处理与故障排查帮助读者在 SeaTunnel Zeta 或 Flink 引擎上快速搭建一套可运行的 PostgreSQL 实时数据管道。支持的引擎与特性PostgreSQL CDC 源连接器插件名Postgres-CDC目前支持以下引擎SeaTunnel Zeta内置引擎Flink从 connector-v2-features 定义的能力矩阵看该连接器当前支持与不支持的能力如下能力支持情况批处理BATCH❌ 不支持流处理STREAMING✅ 支持精确一次Exactly-Once✅ 支持列投影❌ 不支持并行性Parallelism✅ 支持用户自定义拆分✅ 支持其中“精确一次”能力由exactly_once参数开启详见下文并行性由PostgresIncrementalSource实现SupportParallelism接口提供用户自定义拆分则对应snapshot.split.size、table-names-config等表级配置。数据源信息与驱动准备数据源支持的版本驱动JDBC URLMavenPostgreSQL不同依赖版本对应不同驱动类org.postgresql.Driverjdbc:postgresql://localhost:5432/testorg.postgresql:postgresqlPostgreSQL需要操作 GEOMETRY/GEOGRAPHY 类型同上org.postgresql.Driverjdbc:postgresql://localhost:5432/testnet.postgis:postgis-jdbc驱动 JAR 请从 Maven 中央仓库获取本文不再输出外部链接。安装 JDBC 驱动Spark / Flink 引擎将 postgresql 驱动 JAR 放入${SEATUNNEL_HOME}/plugins/目录。SeaTunnel Zeta 引擎将驱动 JAR 放入${SEATUNNEL_HOME}/lib/目录例如cp postgresql-xxx.jar $SEATUNNEL_HOME/lib/连接器源码中硬编码的驱动类名为org.postgresql.Driver见 PostgresSourceConfigFactory.java因此驱动缺失时作业会在启动阶段直接失败。开启 PostgreSQL 逻辑复制必做前置步骤PostgreSQL CDC 依赖逻辑复制必须在服务端完成两步配置第 1 步将wal_level设置为logical编辑postgresql.conf添加wal_level logical修改后需重启 PostgreSQL 服务也可以使用 SQL 命令热修改ALTER SYSTEM SET wal_level TO logical; SELECT pg_reload_conf();第 2 步将监控表的 REPLICA 策略改为 FULL除非require-replica-identity-full显式设置为false否则每张被监控的表都必须执行ALTER TABLE your_table_name REPLICA IDENTITY FULL;REPLICA IDENTITY FULL 会保证 WAL 中携带 UPDATE/DELETE 事件修改前的完整行数据这是 SeaTunnel 正确重建变更前后状态的前提。若未设置作业启动时会报错要求执行该语句见 PostgresCDCIT.java 中 e2e 测试对缺失 FULL 副本标识的校验。数据类型映射PostgreSQL CDC 连接器将 PostgreSQL 类型映射为 SeaTunnel 类型规则如下PostgreSQL 数据类型SeaTunnel 数据类型BOOLBOOLEAN_BOOL数组ARRAYBOOLEANBYTEABYTES_BYTEAARRAYTINYINTINT2、SMALLSERIAL、INT4、SERIALINT_INT2、_INT4ARRAYINTINT8、BIGSERIALBIGINT_INT8ARRAYBIGINTFLOAT4FLOAT_FLOAT4ARRAYFLOATFLOAT8DOUBLE_FLOAT8ARRAYDOUBLENUMERIC列大小 0DECIMAL(列大小, 小数点右侧位数)NUMERIC列大小 0DECIMAL(38, 18)BPCHAR、CHARACTER、VARCHAR、TEXT、GEOMETRY、GEOGRAPHY、JSON、JSONBSTRING_BPCHAR、_CHARACTER、_VARCHAR、_TEXTARRAYSTRINGTIMESTAMPTIMESTAMPTIMETIMEDATEDATE其他数据类型尚不支持映射实现集中在 PostgresTypeUtils.javaNUMERIC的精度判断、数组类型前缀_的处理均可在该文件中找到对应实现。若需要同步 GEOMETRY/GEOGRAPHY 等地理类型请额外引入net.postgis:postgis-jdbc驱动。源选项详解下表汇总了 PostgreSQL CDC 的全部配置项默认值与语义以当前仓库源码为准名称类型必需默认值描述urlString是-JDBC 连接 URL如jdbc:postgresql://localhost:5432/postgres_cdc?loggerLevelOFFusernameString是-连接数据库的用户名passwordString是-连接数据库的密码database-namesList否-需要监控的数据库名称列表table-namesList二选一-需要监控的表使用完整database.schema.table格式如postgres_cdc.inventory.orderstable-patternString二选一-表名正则需匹配完整表名如postgres_cdc\.inventory\..*。与table-names互斥table-names-configList否-表级配置如[{table: db1.schema1.table1, primaryKeys: [key1], snapshotSplitColumn: key2}]。无物理主键的表可通过primaryKeys声明唯一键snapshotSplitColumn必须是唯一键否则被忽略并自动选择拆分列startup.modeEnum否INITIAL启动模式initial先快照后增量、snapshot-only仅快照、有界结束、committed-offset跳过快照从复制槽已提交 LSN 开始要求显式配置slot.name、earliest最早偏移、latest最新偏移stop.modeEnum否NEVER停止模式唯一合法值为never进入增量阶段后持续消费 WALsnapshot.split.sizeInteger否8096快照表拆分大小行数快照阶段按此将表拆成多个 splitsnapshot.fetch.sizeInteger否1024快照读取时每次轮询的最大行数slot.nameString否seatunnelPostgreSQL 逻辑解码槽名称。同一实例多个 CDC 任务必须使用不同的slot.namedecoding.plugin.nameString否pgoutput逻辑解码插件支持decoderbufs、wal2json、wal2json_rds、wal2json_streaming、wal2json_rds_streaming、pgoutputserver-time-zoneString否UTC数据库会话时区未设置时使用ZoneId.systemDefault()connect.timeout.msDuration否30000连接数据库的最大等待时间connect.max-retriesInteger否3建立连接的最大重试次数connection.pool.sizeInteger否20JDBC 连接池大小chunk-key.even-distribution.factor.upper-boundDouble否100块键分布因子上限用于判断表数据是否均匀分布chunk-key.even-distribution.factor.lower-boundDouble否0.05块键分布因子下限sample-sharding.thresholdInteger否1000触发采样分片的估计分片数阈值inverse-sampling.rateInteger否1000采样率的倒数如 1000 表示 1/1000 采样率split.allow-samplingBoolean否true是否允许基于采样的分片策略false 时回退为非均匀分片迭代查询方式enable_concurrent_readBoolean否true快照阶段是否启用基于分片的并发读取false 时跳过分片分析、以单个 split 读整张表适合无索引表exactly_onceBoolean否false快照阶段启用精确一次语义仅initial或snapshot-only可用formatEnum否DEFAULT输出格式可选DEFAULT、COMPATIBLE_DEBEZIUM_JSONrequire-replica-identity-fullBoolean否true是否要求表具有REPLICA IDENTITY FULLfalse 时允许其他副本标识但 UPDATE/DELETE 可能不含修改前状态仅适用于仅追加表如 outbox 模式schema-changes.enabledBoolean否false启用 Schema 演进事件仅支持ADD COLUMN且要求decoding.plugin.name pgoutputschema-changes.includeList否-Schema 演进启用后仅向下游发送列出的事件类型当前支持add.column分组别名update.columns空表示全部schema-changes.excludeList否-不向下游发送列出的 Schema change 事件先 include 后 exclude冲突时 exclude 优先debeziumConfig否-透传给 Debezium 嵌入式引擎的额外属性如心跳配置common-options-否-源插件公共参数详见 源公共选项源码视角参数如何落地startup.mode的五个枚举值在 PostgresSourceOptions.java 中定义stop.mode只允许NEVER。decoding.plugin.name与slot.name的默认值pgoutput/seatunnel定义在 PostgresIncrementalSourceOptions.java此外该文件还定义了schema-name选项与require-replica-identity-full。PostgresSourceConfigFactory.create()会把 SeaTunnel 选项翻译为 Debezium PostgreSQL 连接器属性plugin.name、slot.name、table.include.list、snapshot.mode等并将snapshot-only映射为snapshot.modeinitial_only、committed-offset映射为snapshot.modenever见 PostgresSourceConfigFactory.java。startup.mode committed-offset时若未显式配置slot.namePostgresIncrementalSource.validateStartupOptions()会抛出异常强制要求配置见 PostgresIncrementalSource.java。enable_concurrent_read false时PostgresSourceConfigFactory会将标志传给 ChunkSplitter 以跳过分片分析、整表作为单个 split 读取。快照分片与采样分片策略由 PostgresChunkSplitter.java 实现它继承自 CDC base 模块的AbstractJdbcSourceChunkSplitter其中sampleDataFromColumn负责采样分片的数据抽取。快照分片与采样策略读大数据表的调优点chunk-key.even-distribution.factor的计算公式为(MAX(id) - MIN(id) 1) / 行数若结果落在[lower-bound, upper-bound]默认0.05 ~ 100区间内判定数据均匀分布走均匀分块优化若超出区间且预估分片数近似行数 / 块大小超过sample-sharding.threshold默认 1000则启用基于采样的分片策略采样率为1 / inverse-sampling.rate若split.allow-sampling false无论分片数量多少都回退到非均匀分片迭代查询方式。任务示例从简单同步到高级场景示例一多表实时同步Streaming支持多表读取。该示例将两张 PostgreSQL 表实时同步到同一 PostgreSQL 实例的另一组表带sink_前缀。env { # You can set engine configuration here execution.parallelism 1 job.mode STREAMING checkpoint.interval 5000 read_limit.bytes_per_second7000000 read_limit.rows_per_second400 } source { Postgres-CDC { plugin_output customers_postgres_cdc username postgres password postgres database-names [postgres_cdc] table-names [postgres_cdc.inventory.postgres_cdc_table_1, postgres_cdc.inventory.postgres_cdc_table_2] url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF decoding.plugin.name decoderbufs slot.name seatunnel_postgres_cdc } } transform { } sink { jdbc { plugin_input customers_postgres_cdc url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF driver org.postgresql.Driver username postgres password postgres generate_sink_sql true # You need to configure both database and table database postgres_cdc schema inventory tablePrefix sink_ primary_keys [id] } }关键点说明job.mode STREAMING与checkpoint.interval配合使增量阶段可以周期性记录 LSN 偏移read_limit.bytes_per_second/read_limit.rows_per_second是引擎级的读写限流配置可按需调整table-names使用database.schema.table三段的完整格式源码中PostgresSourceConfigFactory会把三段格式转换为 Debezium 需要的schema.table两段table.include.list。该示例结构即来自仓库 e2e 测试 PostgresCDCIT.java 的真实用例可放心直接复用。示例二ADD COLUMN Schema 演进PostgreSQL 不会把原始ALTER TABLESQL 文本写入逻辑复制流。使用pgoutput插件时SeaTunnel 从 RELATION 消息中检测变更后的表结构在后续行事件之前发出ADD COLUMN事件并更新下游表结构。因此执行 DDL 后该表必须再发生一条行变更连接器才能感知 Schema 变化。如果 RELATION 消息包含ADD COLUMN之外的行 Schema 变更作业会在处理新 Schema 的数据行之前失败从同一 Checkpoint 恢复时会再次遇到该变更只有升级到支持该变更的连接器版本或完成受控的 Schema 迁移并重新启动作业后才能继续处理。source { Postgres-CDC { # ... decoding.plugin.name pgoutput schema-changes.enabled true schema-changes.include [add.column] } }源码层面PostgresIncrementalSource.supports()仅声明支持ADD_COLUMN一种 SchemaChangeType见 PostgresIncrementalSource.java且validateSchemaEvolutionOptions()会在启用 Schema 演进但未使用pgoutput时直接抛异常见 PostgresIncrementalSource.java。RELATION 消息的解析逻辑见 PostgresRelationSchemaChangeResolver.java并有对应的单元测试 PostgresRelationSchemaChangeResolverTest.java。示例三为无主键表自定义主键source { Postgres-CDC { plugin_output customers_postgres_cdc username postgres password postgres database-names [postgres_cdc] table-names [postgres_cdc.inventory.full_types_no_primary_key] url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF decoding.plugin.name decoderbufs exactly_once true slot.name seatunnel_postgres_cdc table-names-config [ { table postgres_cdc.inventory.full_types_no_primary_key primaryKeys [id] } ] } }table-names-config中primaryKeys声明的列会被快照阶段与 WAL 阶段同时用作稳定的行标识配合exactly_once true可获得精确一次保证。示例四配置 Debezium 心跳低流量表保活对于低流量表PostgreSQL 逻辑解码槽的位置只有在 WAL 中发生行变更时才会推进。使用 Debezium 心跳让槽位持续推进便于 checkpoint 定期记录偏移并让复制延迟可观测。心跳表必须提前在 PostgreSQL 服务端创建e2e 测试中即通过CREATE TABLE ... .heartbeat预建见 PostgresCDCIT.java。source { Postgres-CDC { username postgres password postgres database-names [postgres_cdc] schema-names [inventory] table-names [postgres_cdc.inventory.postgres_cdc_table_1] url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF decoding.plugin.name decoderbufs slot.name seatunnel_postgres_cdc debezium { heartbeat.interval.ms 100 heartbeat.action.query INSERT INTO inventory.heartbeat (ts) VALUES (NOW()) } } }示例五仅运行一次性快照BATCH 回填当任务只需执行初始快照并停止不进入 WAL 流式读取时使用startup.mode snapshot-only。该模式适合一次性数据回填。env { parallelism 1 job.mode BATCH checkpoint.interval 5000 } source { Postgres-CDC { username postgres password postgres database-names [postgres_cdc] schema-names [inventory] table-names [postgres_cdc.inventory.postgres_cdc_table_1] url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF decoding.plugin.name decoderbufs slot.name seatunnel_postgres_cdc startup.mode snapshot-only } } sink { Jdbc { url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF driver org.postgresql.Driver username postgres password postgres generate_sink_sql true database postgres_cdc table inventory.sink_postgres_cdc_table_1 primary_keys [id] } }snapshot-only模式下 connector 完全跳过 WAL 流式读取如果快照读取需要独立的复制槽请配置slot.name。从源码看该模式会映射为 Debezium 的snapshot.mode initial_only见 PostgresSourceConfigFactory.java。示例六读取没有主键的表根据源表能提供的保证选择合适的路径仅追加append-only场景源表不会产生 UPDATE/DELETE 事件保持exactly_once false且不声明主键源端会退回到尽力而为的行标识。在没有可用主键的情况下connector 无法安全地应用 UPDATE/DELETE 事件。存在唯一非主键列通过table-names-config.primaryKeys显式声明该列并设置exactly_once true让快照阶段与 WAL 阶段都使用同一配置主键作为稳定的行标识。source { Postgres-CDC { username postgres password postgres database-names [postgres_cdc] schema-names [inventory] table-names [postgres_cdc.inventory.full_types_no_primary_key] url jdbc:postgresql://postgres_cdc_e2e:5432/postgres_cdc?loggerLevelOFF decoding.plugin.name decoderbufs table-names-config [ { table postgres_cdc.inventory.full_types_no_primary_key primaryKeys [id] } ] exactly_once true slot.name seatunnel_postgres_cdc } }没有可用主键时connector 无法安全地应用 UPDATE/DELETE 事件仅在仅追加append-only场景下使用此模式。CDC 元数据字段PostgreSQL CDC 提供以下元数据字段可配合Metadata转换使用字段类型说明databaseSTRING源数据库名称tableSTRING源表名称rowKindSTRING变更类型如 insert、update、deletets_msLONG源事件时间单位毫秒delayLONG事件时间与处理时间之间的延迟单位毫秒示例transform { Metadata { metadata_fields { Database database Table table RowKind rowKind EventTime ts_ms Delay delay } } }常见问题排查PostgreSQL CDC 需要哪些权限CDC 用户需要具备REPLICATION角色以及对监控表的SELECT权限CREATE USER replication_user REPLICATION LOGIN PASSWORD password; GRANT SELECT ON ALL TABLES IN SCHEMA public TO replication_user;同时在postgresql.conf中设置wal_level logical并在pg_hba.conf中添加允许该用户复制连接的条目。支持哪些逻辑解码插件SeaTunnel PostgreSQL CDC 支持pgoutputPostgreSQL 10 起内置、wal2json和decoderbufs默认使用pgoutput。通过decoding.plugin.name参数指定。需要说明的是decoderbufs对应 Debezium 官方插件在使用上要求服务端安装对应的解码插件扩展云数据库如 RDS通常需要选择wal2json_rds/wal2json_rds_streaming变体。完整的六种取值见 PostgresIncrementalSourceOptions.java。SeaTunnel 能从 PostgreSQL 备库读取 CDC 数据吗不能。PostgreSQL 逻辑复制槽必须在主库上创建和消费SeaTunnel 无法直接从备库读取逻辑复制槽需将 CDC 连接器指向主库实例。PostgreSQL CDC 是否支持无主键表默认需要主键。如果表有可作为唯一标识的列可以通过table-names-config中的primaryKeys字段自定义主键参考上文示例三、示例六。复制槽如何管理SeaTunnel 在任务启动时会创建或复用slot.name指定的复制槽。当startup.mode为committed-offset时复制槽必须已存在因为 SeaTunnel 会使用其confirmed_flush_lsn作为启动偏移量。未使用的复制槽会持续占用磁盘上的 WAL 段导致 WAL 持续增长。当 CDC 任务永久下线时应在 PostgreSQL 侧手动删除不再使用的复制槽SELECT * FROM pg_replication_slots; SELECT pg_drop_replication_slot(seatunnel_postgres_cdc);PostgreSQL CDC 为什么会滞后滞后可能由逻辑解码插件处理慢或 WAL sender 负载过高引起。可通过监控pg_replication_slots中的confirmed_flush_lsn漂移情况来排查。确保 CDC 任务持续消费事件并保持 SeaTunnel 与 PostgreSQL 之间的网络低延迟。对于低流量表优先配置 Debezium 心跳见示例四以保持槽位推进。另请参阅若需要一份面向生产的端到端实践指南涵盖全量 增量同步生命周期、2PC sink 配置、Schema 演进与常见故障排查请参阅 CDC 生产实战手册。相关源码与测试入口连接器主体PostgresIncrementalSource.java配置工厂PostgresSourceConfigFactory.java选项定义PostgresIncrementalSourceOptions.java 与 PostgresSourceOptions.javaWAL 读取任务PostgresWalFetchTask.java快照读取任务PostgresSnapshotSplitReadTask.java端到端测试PostgresCDCIT.java赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel Opengauss-CDC 连接器openGauss 快照 WAL 增量实时同步实战指南SeaTunnel Opengauss CDC 连接器openGauss 快照 WAL 增量实时同步实战指南 Apache SeaTunnel 的 Ope数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Opengauss CDC 源连接器从快照到 WAL 增量同步的完整配置与原理指南SeaTunnel Opengauss CDC 源连接器从快照到 WAL 增量同步的完整配置与原理指南 本文以 SeaTunnel 仓库中的 Opengaus数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel MySQL CDC 连接器实战从全量快照到 Binlog 增量同步的完整配置与源码级解析SeaTunnel MySQL CDC 连接器实战从全量快照到 Binlog 增量同步的完整配置与源码级解析 本文基于 SeaTunnel 官方文档 docs数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价