news 2026/9/19 14:22:59

SeaTunnel 实战:PostgreSQL CDC 全量快照 + 增量变更实时写入 Iceberg(字段整形与 Upsert)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SeaTunnel 实战:PostgreSQL CDC 全量快照 + 增量变更实时写入 Iceberg(字段整形与 Upsert)
  • 数据集成
  • ETL
  • 大数据
  • 批处理
  • 流处理
  • 变更数据捕获

【免费下载链接】seatunnel

SeaTunnel 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_source
  • IcebergSink 使用id作为标识字段,以 Upsert 模式应用 CDC 事件,被删除的行也会被正确移除。

本示例属于官方 场景教程(Recipes) 系列,建议先确认你已有的 source 与 sink 组合是否与目标链路一致,再对照envsourcetransformsink四段结构理解参数。

前置条件

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 关键参数

参数在本示例中的值说明
urljdbc:postgresql://postgresql.example.com:5432/salesJDBC 连接串,注意库名直接指向sales(逻辑复制连接实际数据库)
database-names/schema-names["sales"]/["inventory"]需要监控的数据库与 schema
table-names["sales.inventory.customer_orders"]完整database.schema.table格式
startup.modeinitial先同步历史快照,再持续读取增量 WAL
decoding.plugin.namepgoutput逻辑解码插件,源码默认值即pgoutput(见 PostgresIncrementalSourceOptions.java),还支持decoderbufswal2jsonwal2json_rdswal2json_streamingwal2json_rds_streaming
slot.nameseatunnel_sales_orders逻辑复制槽名称,默认seatunnel同一实例上的多个 CDC 任务必须使用不同 slot
debezium.publication.nameseatunnel_sales_orders_pub复用管理员创建的 publication
debezium.publication.autocreate.modedisabled关闭自动创建,保证 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_namerecipe_catalog用户指定的 catalog 名称,默认default
iceberg.catalog.configtype = hadoop,warehouse = ...初始化 Iceberg Catalog 的属性;hadoop类型适用于本地文件/HDFS 等场景,也可换hiveuri = thrift://...)或 S3 Tables REST Catalog
namespace/tablesales_analytics/customer_orders目标库表;不配置table时使用上游表名
iceberg.table.primary-keysid标识行的主键列(逗号分隔)。upsert 模式开启时必须显式配置
iceberg.table.upsert-mode-enabledtrue开启 Upsert 模式,默认false
schema_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_KEYSTABLE_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中应只包含以下两行:

idcustomer_nameamountstatus_namesync_source
1001Alice Zhang150.75PAIDpostgresql_cdc
1003Carol Wang42.00PENDINGpostgresql_cdc

查询 Iceberg 表并核对上述主键集合和每个转换字段:

  • customer_name两侧空白已被trim去除;
  • status已通过upper统一为大写(pendingPENDINGpaidPAID);
  • 每一行都带有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 账号缺少REPLICATIONCONNECT、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.

项目地址:https://gitcode.com/GitHub_Trending/se/seatunnel
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/19 14:17:55

从npm publish到npx使用:命令行工具发布实战与踩坑记录

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/19 14:17:27

Win11本地部署OpenClaw全链路指南:WSL2+Docker+GPU加速实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/19 14:15:47

RIOT 中 LPSXXX 气压传感器驱动:从测试应用到寄存器级实现解析

物联网嵌入式操作系统实时系统 【免费下载链接】RIOT RIOT - The friendly OS for IoT 项目地址: https://gitcode.com/GitHub_Trending/riot/RIOT 点击查看 免费下载 导读 本文围绕 RIOT(The friendly OS for IoT)中 LPSXXX 系列气压传感器…

作者头像 李华
网站建设 2026/9/19 14:14:30

从GPT-3训练算力测算看AI服务器硬件配置与选型逻辑

简介:这是一份聚焦2023年AI服务器市场的行业分析报告,面向算力基础设施、IT硬件及人工智能相关领域的从业者、研究者和投资者。报告基于Counterpoint、IDC等机构数据,指出2022年全球服务器出货量约1380万台、收入1117亿美元,并分析…

作者头像 李华
网站建设 2026/9/19 14:13:20

数字孪生+风景园林:从倾斜摄影到积温驱动的季相演算

简介:这是一份PDF学术资料,围绕数字孪生技术在风景园林设计中的应用展开,适合风景园林设计师、研究人员以及智慧城市相关从业者阅读。内容从数字孪生技术概述切入,重点阐述其与LIM风景园林信息模型的融合路径,强调实时…

作者头像 李华