SeaTunnel 实战:PostgreSQL CDC 全量快照 + 增量变更实时写入 Iceberg(字段整形与 Upsert)
数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载本教程是 SeaTunnel Zeta 引擎下的一条完整实时同步链路先通过Postgres-CDC连接器读取 PostgreSQL 全量快照并持续捕获增量变更再使用SqlTransform 对字段做清洗、值域归一化并补充来源标记最后由IcebergSink 以主键 Upsert 模式在数据湖中维护每一行的最新状态。读完本文你将掌握 PostgreSQL 逻辑复制的环境准备、只读 CDC 账号与 publication 的最小权限设计、HOCON 任务配置的每个关键参数以及 Iceberg Delta Writer 处理 CDC 事件快照/更新/新增/删除的底层原理。场景与完整链路本场景示例使用 PostgreSQL 14 与pgoutput逻辑解码插件数据源为sales.inventory.customer_orders表覆盖快照、更新、新增和删除四类事件。完整链路如下Postgres-CDC通过 PostgreSQL 逻辑复制读取sales.inventory.customer_orders的存量数据与增量 WAL 变更SqlTransform 清理客户名称trim、统一状态值upper并填充来源标记sync_sourceIcebergSink 使用id作为标识字段以 Upsert 模式应用 CDC 事件被删除的行也会被正确移除。本示例属于官方 场景教程Recipes 系列建议先确认你已有的 source 与 sink 组合是否与目标链路一致再对照env、source、transform、sink四段结构理解参数。前置条件1. 先运行第一个本地任务请先使用 Zeta 引擎完成运行第一个任务确保本机SEATUNNEL_HOME环境与启动脚本可用。2. 安装所需 Connector在${SEATUNNEL_HOME}下的config/plugin_config中声明本次任务需要的连接器--seatunnel-connectors-- connector-cdc-postgres connector-iceberg --end--然后执行安装脚本并确认插件已就位cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | grep -E connector-(cdc-postgres|iceberg)3. 放置 PostgreSQL JDBC 驱动Zeta 引擎需要把兼容版本的 PostgreSQL JDBC 驱动放入${SEATUNNEL_HOME}/libls ${SEATUNNEL_HOME}/lib | grep postgresql-这是 Zeta 引擎与 Spark/Flink 引擎的区别Spark/Flink 场景驱动放在${SEATUNNEL_HOME}/plugins/而 Zeta 场景统一放在${SEATUNNEL_HOME}/lib/详见 PostgreSQL CDC 连接器文档。4. 在主库启用逻辑复制在 PostgreSQLpostgresql.conf中修改以下参数修改wal_level后必须重启 PostgreSQLwal_level logical max_replication_slots 10 max_wal_senders 10提示如果不便重启也可使用ALTER SYSTEM SET wal_level TO logical;后再SELECT pg_reload_conf();但wal_level属于需要重启才能生效的参数postmaster级务必确认实际生效值。确认实际生效值SHOW wal_level; SHOW max_replication_slots; SHOW max_wal_senders;5. 创建源数据库与只读 CDC 账号使用 PostgreSQL 管理员账号创建源数据库和专用 CDC 用户并按实际环境替换数据库、用户名和密码如果数据库或用户已存在直接复用并跳过对应的CREATE语句CREATE DATABASE sales;PostgreSQL 要求在事务外执行CREATE DATABASE。创建完成后再执行CREATE ROLE seatunnel_cdc WITH REPLICATION LOGIN PASSWORD change_me; GRANT CONNECT ON DATABASE sales TO seatunnel_cdc;在pg_hba.conf中添加匹配规则并把示例 CIDR 替换为 SeaTunnel 节点所在的实际网段host sales seatunnel_cdc 192.0.2.0/24 scram-sha-256这里有一点容易踩坑逻辑复制连接的是实际数据库因此这一行填写的是数据库名sales而不是物理复制连接streaming replication所使用的replication关键字。保存文件后重新加载 PostgreSQL 配置SELECT pg_reload_conf();6. 准备 Iceberg Warehouse选择所有 SeaTunnel Worker 都能访问且可写的空Iceberg warehouse。只有所有 Worker 共享同一文件系统时才适合使用本地file://路径分布式集群应使用 HDFS、S3 或其他共享 catalog 存储。多 Worker 集群使用节点本地路径会导致不同 Worker 写入的文件彼此不可见任务会在 checkpoint 提交或后续读取时出现异常。准备源数据使用管理员账号连接sales数据库在 PostgreSQL 14 中创建源 schema 和表然后为只读 CDC 用户授权CREATE SCHEMA inventory; CREATE TABLE inventory.customer_orders ( id BIGINT PRIMARY KEY, customer_name VARCHAR(64) NOT NULL, amount NUMERIC(10, 2) NOT NULL, status VARCHAR(16) NOT NULL, updated_at TIMESTAMP NOT NULL ); ALTER TABLE inventory.customer_orders REPLICA IDENTITY FULL; INSERT INTO inventory.customer_orders VALUES (1001, Alice Zhang , 120.50, pending, 2026-07-18 09:00:00), (1002, Bob Li, 80.00, paid, 2026-07-18 09:05:00); GRANT USAGE ON SCHEMA inventory TO seatunnel_cdc; GRANT SELECT ON TABLE inventory.customer_orders TO seatunnel_cdc; CREATE PUBLICATION seatunnel_sales_orders_pub FOR TABLE inventory.customer_orders;几点关键说明REPLICA IDENTITY FULLSeaTunnel 的 PostgreSQL CDC 连接器默认要求源表设置REPLICA IDENTITY FULL对应源码中的require-replica-identity-full选项默认true见 PostgresIncrementalSourceOptions.java否则 UPDATE/DELETE 事件可能不包含变更前的完整行状态默认安全检查会直接报错。Publication 由管理员创建因此 CDC 用户可以保持只读权限CONNECT schemaUSAGE 表SELECT不需要在源库上拥有写权限。复用而非自动创建 Publication下面的任务配置通过debezium.publication.autocreate.mode disabled复用上面已创建的 publication而不是尝试为数据库中的所有表自动创建。初始快照进入 Iceberg 后使用源业务的写入账号或管理员账号执行以下变更不要使用只读 CDC 账号用于验证后续的更新、新增、删除三类增量事件UPDATE inventory.customer_orders SET amount 150.75, status paid, updated_at 2026-07-18 10:00:00 WHERE id 1001; INSERT INTO inventory.customer_orders VALUES (1003, Carol Wang , 42.00, pending, 2026-07-18 10:05:00); DELETE FROM inventory.customer_orders WHERE id 1002;完整任务配置下面的 HOCON 实现了完整链路。请按实际环境替换示例主机名、账号、slot 名称和 warehouse 路径。同一个 PostgreSQL 实例上的并发 CDC 任务必须使用不同的slot.name并确保publication.name与上面创建的 publication 一致。env { parallelism 1 job.mode STREAMING checkpoint.interval 3000 } source { Postgres-CDC { plugin_output postgres_orders_raw url jdbc:postgresql://postgresql.example.com:5432/sales username seatunnel_cdc password change_me database-names [sales] schema-names [inventory] table-names [sales.inventory.customer_orders] startup.mode initial decoding.plugin.name pgoutput slot.name ${slot_name} debezium { publication.name seatunnel_sales_orders_pub publication.autocreate.mode disabled } } } transform { Sql { plugin_input postgres_orders_raw plugin_output iceberg_customer_orders query select id, trim(customer_name) as customer_name, amount, upper(status) as status_name, updated_at, postgresql_cdc as sync_source from dual } } sink { Iceberg { plugin_input iceberg_customer_orders catalog_name recipe_catalog iceberg.catalog.config { type hadoop warehouse ${warehouse} } namespace sales_analytics table customer_orders iceberg.table.primary-keys id iceberg.table.upsert-mode-enabled true schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }把配置保存为${SEATUNNEL_HOME}/config/postgresql-cdc-to-iceberg.conf将${slot_name}和${warehouse}替换为带引号的实际值然后运行cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/postgresql-cdc-to-iceberg.conf -m local例如slot.name seatunnel_sales_orders warehouse file:///tmp/seatunnel/iceberg/postgres-cdc-recipe/配置参数深度解析Source 段Postgres-CDC 关键参数参数在本示例中的值说明urljdbc:postgresql://postgresql.example.com:5432/salesJDBC 连接串注意库名直接指向sales逻辑复制连接实际数据库database-names/schema-names[sales]/[inventory]需要监控的数据库与 schematable-names[sales.inventory.customer_orders]完整database.schema.table格式startup.modeinitial先同步历史快照再持续读取增量 WALdecoding.plugin.namepgoutput逻辑解码插件源码默认值即pgoutput见 PostgresIncrementalSourceOptions.java还支持decoderbufs、wal2json、wal2json_rds、wal2json_streaming、wal2json_rds_streamingslot.nameseatunnel_sales_orders逻辑复制槽名称默认seatunnel同一实例上的多个 CDC 任务必须使用不同 slotdebezium.publication.nameseatunnel_sales_orders_pub复用管理员创建的 publicationdebezium.publication.autocreate.modedisabled关闭自动创建保证 CDC 账号只读startup.mode还有多种可选值按需选择详见 PostgreSQL CDC 连接器文档snapshot-only仅同步启动时的历史数据之后以有界任务结束不进入 WAL 流读取committed-offset跳过快照从配置的复制槽已提交 LSN 开始读取要求显式配置slot.nameslot 不存在或无已提交 LSN 时启动失败earliest/latest分别从最早偏移量、最新偏移量启动。快照阶段的拆分与读取也支持细粒度调优snapshot.split.size默认 8096 行、snapshot.fetch.size默认 1024、chunk-key.even-distribution.factor上下限默认 100 / 0.05、sample-sharding.threshold默认 1000 分片等适合大表分片扫描场景。Transform 段Sql 字段整形SqlTransform 使用内存 SQL 引擎通过plugin_input/plugin_output衔接上下游表名query中表名必须与plugin_input一致from dual表示对单行事件做表达式计算。本示例在一条 SQL 里完成了三类加工select id, trim(customer_name) as customer_name, -- 清理客户名称两侧空白 amount, upper(status) as status_name, -- 状态值统一为大写 updated_at, postgresql_cdc as sync_source -- 填充来源标记常量 from dual注意SQL Transform 支持函数与条件过滤但不支持多源表 JOIN 和聚合等复杂 SQL嵌套结构查询中不能出现table_name。详见 Sql Transform 文档。Sink 段Iceberg 参数解析参数在本示例中的值说明catalog_namerecipe_catalog用户指定的 catalog 名称默认defaulticeberg.catalog.configtype hadoop,warehouse ...初始化 Iceberg Catalog 的属性hadoop类型适用于本地文件/HDFS 等场景也可换hiveuri thrift://...或 S3 Tables REST Catalognamespace/tablesales_analytics/customer_orders目标库表不配置table时使用上游表名iceberg.table.primary-keysid标识行的主键列逗号分隔。upsert 模式开启时必须显式配置iceberg.table.upsert-mode-enabledtrue开启 Upsert 模式默认falseschema_save_modeCREATE_SCHEMA_WHEN_NOT_EXISTschema 不存在时自动创建这也是源码默认值data_save_modeAPPEND_DATA数据写入方式默认值即APPEND_DATA可选CUSTOM_PROCESSING配合custom_sql在写入前执行 delete以上选项的默认值与约束在 IcebergSinkOptions.java 中均有明确定义iceberg.table.upsert-mode-enabled默认falseiceberg.table.primary-keys默认无值且当 upsert 开启时必须非空——Sink 不会再自动继承 Source 表的主键见该文件TABLE_PRIMARY_KEYS与TABLE_UPSERT_MODE_ENABLED_PROP的注释。此外还有iceberg.table.schema-evolution-enabled默认false开启后同步过程中支持 schema 变更、iceberg.table.write-props透传给写入器如write.format.default parquet、write.target-file-size-bytes优先级最高等可选配置。Upsert 的底层实现Delta Writer从源码结构看Iceberg Sink 的写入器选择逻辑在 IcebergWriterFactory.java 中当identifierFieldIds由iceberg.table.primary-keys解析而来为空且未开启 upsert 模式时使用普通UnpartitionedWriter/PartitionedAppendWriter仅做追加写只要配置了主键或开启了 upsert 模式就改用UnpartitionedDeltaWriter/PartitionedDeltaWriter。Delta Writer 依据主键区分 CDC 事件行携带后像after的行按主键执行 insert/update携带前像before的删除事件则生成对应的删除操作从而在 Iceberg 表上维护“每行一个最新状态”的语义。这正是本示例能够在 Iceberg 中正确反映UPDATEid1001 金额与状态变更、INSERTid1003与DELETEid1002 被移除三种事件的底层原因Upsert 模式要求显式主键主键必须对应源表稳定的唯一键否则无法将增量事件与目标行正确关联。运行与验证提交配置并启动任务后增量 SQL 提交且下一次 Iceberg checkpoint 成功sales_analytics.customer_orders中应只包含以下两行idcustomer_nameamountstatus_namesync_source1001Alice Zhang150.75PAIDpostgresql_cdc1003Carol Wang42.00PENDINGpostgresql_cdc查询 Iceberg 表并核对上述主键集合和每个转换字段customer_name两侧空白已被trim去除status已通过upper统一为大写pending→PENDING、paid→PAID每一行都带有sync_source postgresql_cdc来源标记被删除的1002不得继续存在——这是对 Upsert Delta Writer 删除语义最直接的验证点。运行检查与运维要点持续监控pg_replication_slots停止消费的 slot 可能无限保留 WAL导致磁盘膨胀与下游延迟只有对应 CDC 任务永久下线后才能删除 slot活动任务之间禁止共用 sloticeberg.table.primary-keys必须对应稳定的源表主键upsert 模式要求显式配置该参数确认 checkpoint 持续成功Iceberg 变更会在成功提交后可见本示例checkpoint.interval 3000可按吞吐需求调整多 Worker 集群不要使用节点本地 warehouse除非该路径实际由共享存储承载。常见问题排查wal_level不是logical或者修改配置后没有重启 PostgreSQLSHOW wal_level确认CDC 账号缺少REPLICATION、CONNECT、schemaUSAGE或表SELECT权限配置的 replication slot 已被另一个任务使用改用唯一slot.name默认安全检查开启时源表没有设置REPLICA IDENTITY FULLIceberg warehouse 不可写或并非所有 SeaTunnel Worker 都可见开启 upsert 模式却没有显式设置iceberg.table.primary-keys。相关文档PostgreSQL CDC Source 连接器文档Iceberg Sink 连接器文档Sql Transform 文档场景教程总览赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel 实战PostgreSQL CDC 实时同步到 Iceberg字段整形与 Upsert 完整指南SeaTunnel 实战PostgreSQL CDC 实时同步到 Iceberg字段整形与 Upsert 完整指南 本篇技术指南基于 Apache Sea数据集成ETL大数据批处理流处理变更数据捕获Flink CDC 实战PostgreSQL CDC Connector 全量快照与增量变更捕获完整指南Flink CDC 实战PostgreSQL CDC Connector 全量快照与增量变更捕获完整指南 本文基于 Apache Flink CDC 开源仓库后端数据集成大数据流处理变更数据捕获数据同步SeaTunnel Opengauss-CDC 连接器openGauss 快照 WAL 增量实时同步实战指南SeaTunnel Opengauss CDC 连接器openGauss 快照 WAL 增量实时同步实战指南 Apache SeaTunnel 的 Ope数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →