- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
本教程是 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_source;IcebergSink 使用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}/lib:
ls "${SEATUNNEL_HOME}/lib" | grep 'postgresql-'这是 Zeta 引擎与 Spark/Flink 引擎的区别:Spark/Flink 场景驱动放在${SEATUNNEL_HOME}/plugins/,而 Zeta 场景统一放在${SEATUNNEL_HOME}/lib/,详见 PostgreSQL CDC 连接器文档。
4. 在主库启用逻辑复制
在 PostgreSQLpostgresql.conf中修改以下参数,修改wal_level后必须重启 PostgreSQL:
wal_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 FULL:SeaTunnel 的 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 关键参数
| 参数 | 在本示例中的值 | 说明 |
|---|---|---|
url | jdbc:postgresql://postgresql.example.com:5432/sales | JDBC 连接串,注意库名直接指向sales(逻辑复制连接实际数据库) |
database-names/schema-names | ["sales"]/["inventory"] | 需要监控的数据库与 schema |
table-names | ["sales.inventory.customer_orders"] | 完整database.schema.table格式 |
startup.mode | initial | 先同步历史快照,再持续读取增量 WAL |
decoding.plugin.name | pgoutput | 逻辑解码插件,源码默认值即pgoutput(见 PostgresIncrementalSourceOptions.java),还支持decoderbufs、wal2json、wal2json_rds、wal2json_streaming、wal2json_rds_streaming |
slot.name | seatunnel_sales_orders | 逻辑复制槽名称,默认seatunnel;同一实例上的多个 CDC 任务必须使用不同 slot |
debezium.publication.name | seatunnel_sales_orders_pub | 复用管理员创建的 publication |
debezium.publication.autocreate.mode | disabled | 关闭自动创建,保证 CDC 账号只读 |
startup.mode还有多种可选值,按需选择(详见 PostgreSQL CDC 连接器文档):
snapshot-only:仅同步启动时的历史数据,之后以有界任务结束,不进入 WAL 流读取;committed-offset:跳过快照,从配置的复制槽已提交 LSN 开始读取,要求显式配置slot.name,slot 不存在或无已提交 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_name | recipe_catalog | 用户指定的 catalog 名称,默认default |
iceberg.catalog.config | type = hadoop,warehouse = ... | 初始化 Iceberg Catalog 的属性;hadoop类型适用于本地文件/HDFS 等场景,也可换hive(uri = thrift://...)或 S3 Tables REST Catalog |
namespace/table | sales_analytics/customer_orders | 目标库表;不配置table时使用上游表名 |
iceberg.table.primary-keys | id | 标识行的主键列(逗号分隔)。upsert 模式开启时必须显式配置 |
iceberg.table.upsert-mode-enabled | true | 开启 Upsert 模式,默认false |
schema_save_mode | CREATE_SCHEMA_WHEN_NOT_EXIST | schema 不存在时自动创建,这也是源码默认值 |
data_save_mode | APPEND_DATA | 数据写入方式,默认值即APPEND_DATA;可选CUSTOM_PROCESSING配合custom_sql在写入前执行 delete |
以上选项的默认值与约束在 IcebergSinkOptions.java 中均有明确定义:iceberg.table.upsert-mode-enabled默认false;iceberg.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 中正确反映UPDATE(id=1001 金额与状态变更)、INSERT(id=1003)与DELETE(id=1002 被移除)三种事件的底层原因:Upsert 模式要求显式主键,主键必须对应源表稳定的唯一键,否则无法将增量事件与目标行正确关联。
运行与验证
提交配置并启动任务后,增量 SQL 提交且下一次 Iceberg checkpoint 成功,sales_analytics.customer_orders中应只包含以下两行:
| id | customer_name | amount | status_name | sync_source |
|---|---|---|---|---|
| 1001 | Alice Zhang | 150.75 | PAID | postgresql_cdc |
| 1003 | Carol Wang | 42.00 | PENDING | postgresql_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;活动任务之间禁止共用 slot;
iceberg.table.primary-keys必须对应稳定的源表主键,upsert 模式要求显式配置该参数;- 确认 checkpoint 持续成功,Iceberg 变更会在成功提交后可见(本示例
checkpoint.interval = 3000,可按吞吐需求调整); - 多 Worker 集群不要使用节点本地 warehouse,除非该路径实际由共享存储承载。
常见问题排查
wal_level不是logical,或者修改配置后没有重启 PostgreSQL(SHOW wal_level确认);- CDC 账号缺少
REPLICATION、CONNECT、schemaUSAGE或表SELECT权限; - 配置的 replication slot 已被另一个任务使用(改用唯一
slot.name); - 默认安全检查开启时,源表没有设置
REPLICA IDENTITY FULL; - Iceberg warehouse 不可写,或并非所有 SeaTunnel Worker 都可见;
- 开启 upsert 模式,却没有显式设置
iceberg.table.primary-keys。
相关文档
- PostgreSQL CDC Source 连接器文档
- Iceberg Sink 连接器文档
- Sql Transform 文档
- 场景教程总览
- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
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),仅供参考