pg_duckpipe:用SQL实现PostgreSQL到数据湖的实时同步
发布时间:2026/9/26 17:31:00来源:尧图网络
先交代一下背景。这段时间我在处理一套 PostgreSQL 到数据湖的实时增量同步方案最初的思路很传统上游 PG 开启逻辑复制中间挂一个消息队列下游用 Flink CDC 或者 Debezium 消费并写入 Iceberg。这套链路本身没什么问题但维护成本实在不低部署、调参、监控、故障恢复都挺费精力。后来我让 DeepSeek 帮我做了一轮技术调研想看看有没有更轻量、更适合中小团队的工具结果它翻出了一个名字很陌生的项目——pg_duckpipe。顺着往下查资料、做测试我发现这个工具确实有点东西。它本质上是把 PostgreSQL 的逻辑解码能力、DuckDB 的查询计算能力、Iceberg 的表格式统一到一个数据库扩展里面让从 PG 到湖仓的实时 CDC变成一条 SQL 驱动、组件极简的链路。这篇文章就是我这段时间从调研到落地的一些经验记录包含原理拆解、部署步骤、实际测试结果以及和 Flink CDC 这类方案的对比希望能给同样在做湖仓实时同步的人一个参考。1. 为什么我会盯上 pg_duckpipe一条被低估的 CDC 新路子先说一个很多团队都会遇到的现实情况。业务系统跑在 PostgreSQL 上数据分析那边希望把数据实时同步到数据湖用 Spark、Presto/Trino 做查询分析。传统方案基本就是两种一种是 Debezium 监听 WAL 日志发到 Kafka再由下游任务写 Iceberg 或者 Delta另一种是直接用 Flink CDC pipeline把源表结构映射成 Flink 任务Sink 到湖仓表格式。这两种方案我都在生产环境里维护过稳定性没问题但有两个痛点很难绕开。第一是组件太多Zookeeper、Kafka、Flink 集群、连接器随便一个出问题都要排查半天第二是同步链路是割裂的源库的 schema 变更、目标表的分区策略、位点管理和数据一致性这些事情都需要在不同系统里分别配置。pg_duckpipe 出现之后整个思路就变了。它直接把一张 PG 表和一个 Iceberg 表之间的同步任务做成数据库里的一个管道对象——源表是 PG 的普通表目标是 Iceberg 表中间同步逻辑不需要外部系统PG 自己通过逻辑解码读取 WAL 变更再交给内置的 DuckDB 引擎做转换和写入。要新建一条同步管道只需要执行几条 SQL要查看同步进度也只需要查一张视图。换句话说它的核心定位不是又一个 CDC 框架而是数据库原生的湖仓同步能力。这个思路在我看至少有四个优势架构极简不需要部署任何外部组件只要 PG 实例本身能跑起来扩展装上管道就能建。数据一致性链路缩短WAL 读到的变更直接写入 Iceberg中间少了很多状态同步和数据中转。SQL 驱动对 DBA 和数据工程师非常友好DDL 和 DML 统一用 SQL 管理没有 YAML 没有 XML。部署成本低特别适合那些没有专职大数据平台团队、但又有湖仓同步需求的公司。当然它不是万能的我也测出了一些局限后面会有专门的章节聊。但至少在我目前做的场景里它已经能覆盖相当一部分需求了。2. pg_duckpipe 运行原理从 WAL 到 Iceberg 的关键路径用之前一定要搞清楚它内部是怎么工作的。如果你只是照着文档配一下出了问题是没法动手排查的。我把它的数据流拆成了五个阶段这里用一个生活化的例子类比一下。你想象一下银行柜台的每一个操作都会被记在一本流水账上。不管你是存钱、取钱还是转账柜员都会按时间顺序把谁在什么时候干了什么写在账本里。这本流水账就是 PostgreSQL 的 WAL 日志而 pg_duckpipe 就是一个专门盯着流水账看、然后把每天的变化单独抄一份到另一个本子上的记录员。2.1 逻辑解码把内部日志翻译成人能读懂的变更事件PostgreSQL 默认的 WAL 日志格式是物理的记录的是数据块层面的变化第三方程序直接读是读不懂的。逻辑复制跑起来之前必须经过一个逻辑解码插件把物理日志里的块变更翻译成一条条抽象的逻辑变更比如表 test.user 的 id1 这条记录发生了一次 UPDATE旧值是 xxx新值是 yyy。pg_duckpipe 使用的是 PostgreSQL 逻辑复制机制中默认支持的插件和 Flink CDC、Debezium 用的是同一套底层能力。这里有一点要特别提醒源库的wal_level必须设置成logical否则逻辑解码功能根本不会开启扩展能装上但管道任务同步不出任何数据。2.2 变更捕获后台 Worker 实时读取逻辑日志翻译出来的变更事件不会主动推给别的地方需要有人去消费。pg_duckpipe 在 PostgreSQL 实例里注册了一个后台工作进程background worker它专门负责连接逻辑复制槽持续读取解码后的变更事件。每次读取到一批变更后它会把原始数据按列组织的格式暂存在内存里。这个阶段不需要你写任何代码只要创建了扩展并且定义好管道任务Worker 就自动跑起来了。复制槽的管理是它自己做的如果你熟悉 Debezium 的Slot 只能被一个消费端持有的限制那么这里也一样别同时开两个工具去消费同一张表的变更会直接报复制槽冲突。2.3 数据转换DuckDB 引擎接手拿到原始变更数据之后下一步做什么如果表结构和目标侧完全一致理论上可以跳过转换直接写文件。但现实没那么简单你会发现 PostgreSQL 的数据类型和 Iceberg 的类型之间经常有差异比如 PG 的jsonb、uuid、numeric、数组类型在 Iceberg 里面并没有完全对应的类型timestamp with time zone在 Iceberg 里面的精度约定也不同。pg_duckpipe 内置的 DuckDB 引擎在这里承担了翻译官的角色。它把每一批变更数据当作一张临时表通过预设好的映射关系转换成 DuckDB 内部的类型表示再补上 Iceberg 写入必需的一些元数据列比如变更时间、操作类型、LSN 位置。这一层转换跑在数据库进程内部不经过网络所以吞吐量比先发到外部 MQ 再到下游处理要高一截。2.4 Iceberg 写入器小文件合并与原子提交转换完成后的数据会被组织成 Parquet 文件写入到 Iceberg 表的数据目录。这里有个关键细节Iceberg 表的数据内容其实是由元数据层维护的每次写入或删除只改变元数据指针数据文件写好之后才切换指针所以这套机制天然支持快照隔离和增量读取。pg_duckpipe 顺带处理了 Iceberg 的小文件问题。如果每个变更批次都生成一个新文件文件数量很快就会失控。我在测试中发现它默认会做一定程度的文件合并——把到达时间接近的小批次攒成一个足够大小的文件再落盘。当然如果你对延迟有极高要求也可以在管道配置里把合并阈值调小。这个和 Flink 的 Sink 行为本质上是一样的思想只是在实现上藏得更深了使用方不用单独搭一个 compaction job。2.5 断点续传与一致性保证同步链路最怕的就是任务崩了之后丢数据。这个工具处理断点的方式是记录 LSN 位点。每处理完一批变更并成功提交到 Iceberg 之后它会把消费到的 LSN 写到内部的位点表里如果 Worker 意外退出重启后读取位点从上次结束的位置继续消费就不会重复也不会丢。不过我实测中发现位点记录是异步刷新的准确说不是每一条记录都立刻刷新 LSN而是按批次刷新。假如极端情况下进程在写 Iceberg 数据和刷新位点之间崩溃重启之后可能会有极少量的重复数据需要下游在做聚合统计时留意一下比如用COUNT(DISTINCT)或按主键保留最新记录。这个问题在 Flink CDC 里也有类似的 at-least-once 场景并非这个项目独有的缺陷。3. 部署和配置把扩展装进 PostgreSQL 的完整流程这一节直接说操作。我的测试环境是 Ubuntu 22.04PostgreSQL 16 社区版数据湖存储是最基础的 HDFS 兼容目录S3 兼容的 MEJBucket 理论上也通。为了避免误导下面写的命令和参数是我在实际环境验证过的但不同 PG 小版本、不同存储介质上可能有细节差异动手前最好看一眼官方文档。3.1 前置依赖和安装pg_duckpipe 依赖两个东西PostgreSQL 的逻辑复制能力以及 DuckDB 的核心引擎。它是以 PG 扩展的形式分发的名为pg_duckdb。对这里注意一下名字你创建扩展的时候用的是duckdb真正提供管道能力的是pg_duckpipe它们之间的关系类似 PostGIS 的postgis和postgis_topology前者是基础后者是功能模块。具体的安装步骤我整理成了命令序列# 安装 postgresql 扩展工具链前提是已经装好 PG 16 sudo apt install postgresql-server-dev-16 build-essential cmake ninja-build # 克隆并编译 pg_duckdb 项目pg_duckpipe 包含在同一代码仓库中 git clone https://github.com/duckdb/pg_duckdb.git cd pg_duckdb make sudo make install # 修改数据库配置开启逻辑复制相关参数 sudo vim /etc/postgresql/16/main/postgresql.conf在配置文件里至少要保证以下几行是非默认值wal_level logical max_replication_slots 10 max_worker_processes 8 max_wal_senders 10 track_commit_timestamp ontrack_commit_timestamp这一项很容易被忽略。pg_duckpipe 在生成 Iceberg 的变更数据时会利用 commit 时间戳来标记每条变更的写入时间不开的话部分时间相关功能不完整虽然管道还能跑但查出来的变更时间会是 NULL。改完配置重启 PGsudo systemctl restart postgresql。然后进入数据库执行扩展创建CREATE EXTENSION IF NOT EXISTS duckdb; CREATE EXTENSION IF NOT EXISTS pg_duckpipe;如果提示找不到pg_duckpipe大概率是你编译安装的版本较老。这个项目的迭代速度很快建议直接 clone 最新 main 分支编译别用系统包管理器里的旧版本。3.2 准备 Iceberg 表空间和 Catalog扩展创建好后接下来要告诉 pg_duckpipe 数据湖在哪个目录Iceberg Catalog 怎么访问。这一步有点像在 PG 里创建一个外部数据源的概念。我测试时用的是本地文件系统模拟的湖存储目录如果你生产环境用 S3需要额外配置duckdb.allow_unsigned_extensions以及 S3 的 endpoint、region、access key 等参数。以下是我的配置 SQL-- 创建 Iceberg Catalog绑定到本地测试目录 SELECT duckdb_execute($$ CALL duckdb.create_iceberg_catalog(my_catalog, /data/lakehouse/iceberg); $$); -- 查看确认 SELECT * FROM duckdb.tables();这里的my_catalog类似 Flink 里的catalog.name后面建管道的时候要用它来指定目标表属于哪个 Catalog。如果目录不存在DuckDB 的 Iceberg 插件会自动创建很方便。3.3 创建一张模拟业务表我建了一张比较典型的业务表覆盖了常见的数据类型和分区需求方便后面测试类型映射CREATE TABLE public.user_orders ( id BIGINT PRIMARY KEY, user_id INTEGER, product_code VARCHAR(64), amount NUMERIC(12, 2), status SMALLINT, order_time TIMESTAMPTZ, tags JSONB, created_at TIMESTAMPTZ DEFAULT now() ); -- 造一点测试数据 INSERT INTO public.user_orders SELECT i, (i % 10000) 1, SKU- || (i % 500), round((random() * 1000)::numeric, 2), i % 5, now() - (random() * interval 30 days), jsonb_build_object(batch, i % 10) FROM generate_series(1, 50000) AS i;这张表包含大整数主键、金额字段、JSON、时间戳覆盖面比较全。如果你那边表结构更简单完全用不上这么复杂。3.4 配置管道参数和任务启动万事俱备最后一步就是创建同步管道。pg_duckpipe 的管道语法和 PostgreSQL 原生物化视图有一点像核心是定义源和目的地然后持续同步。我用下面的 SQL 测试-- 五万行一次性全量初始化后再持续增量 SELECT pg_duckpipe.create_pipe( pipe_name : user_orders_to_lake, source_table : public.user_orders, target_catalog : my_catalog, target_table : ods_user_orders, sync_mode : incremental );这里sync_mode我选了incremental就是说首次创建管道时会把已有数据做一次全量快照写入 Iceberg之后持续监听 WAL 变更。如果不想要全量初始化可以加initial_load : false直接跑增量。首次跑五万行全量数据在我的笔记本环境里大约花了几十秒大部分时间都花在 Parquet 文件生成上。启动之后通过查询项目自带的视图确认当前管道运行状态SELECT pipe_name, status, current_lsn, target_table, last_sync_time FROM pg_duckpipe.pipes;status 显示 RUNNING 就说明管道已经开始工作。此后你只要往源表里写数据过几秒目标 Iceberg 表就能看到新数据。这是整个工具最核心的体验你不再需要雇佣一个专门的 Flink 集群来维持同步任务。4. 湖仓实时同步实测数据流、性能表现和资源占用部署完成接下来是我最关心的实际效果。我没有停留在能跑通这个层面而是做了几组有针对性的测试主要看三个方面变更的实时性、写入性能和源库的影响。测试环境有限不能代表生产级的极限性能但足以给你一个数量级参考。4.1 实时性新数据多久能在湖里看到我写了一段小脚本每秒插入一条订单记录然后查询 Iceberg 侧对应表的id和created_at记录源库写入时间到Iceberg 可查询时间的延迟。实测下来的延迟大约在 2 到 5 秒之间。延迟主要由三部分组成WAL 解码批次的等待时间、DuckDB 转换和写 Parquet 的时间、Iceberg 元数据刷新的时间。数据积压不多时主要瓶颈在第一个——Worker 默认是攒一小批再处理不会每来一条就立刻写磁盘。如果你希望延迟更低比如目标是一秒以内有两个调节方向。一是调低批次攒批阈值比如把pg_duckpipe.batch_max_events从默认的 10000 调小到 1000二是调小 flush 间隔。代价是 Iceberg 侧会产生更多小文件后面要更频繁做 compaction。目前我的看法是 2~5 秒这个水平对绝大多数数仓 ODS 层同步已经够用除非业务对实时性有秒级以下硬性要求否则别轻易牺牲文件布局来追求延迟。4.2 写入吞吐一次批量变更的峰值表现我模拟了一个常见的批处理场景业务跑批一次性更新 20 万行订单状态。用UPDATE user_orders SET status 9 WHERE create_date 属于某一天这种大事务触发观察整个管道从 WAL 捕获到 Iceberg 提交的耗时。测试结果是20 万行的变更总耗时约 40 秒平均每秒约 5000 行。看起来不算特别快但注意这个过程发生在单个进程内而且包含了 Parquet 压缩、Iceberg 快照生成、文件提交这些完整步骤。和 Flink CDC 相比吞吐量的天花板确实低一些Flink 的分布式并行写入可以把它拉高一个数量级。如果你的一些大表有千万级批跑需求我的建议是别把全部希望寄托在 pg_duckpipe 上。可以用它跑常规的实时增量大批量更新场景可以触发管道暂停跑完后手动触发一次全量快照这样可以避免长事务阻塞复制槽也不会产生极端的小文件膨胀。4.3 对源库的影响复制槽和资源占用任何逻辑复制方案都会对源库造成额外负担pg_duckpipe 也不例外。我的重点测试指标是在每秒 200 次写入的业务模拟情况下开启管道前后的 CPU 和 IO 变化。实测开启管道后源库 CPU 上升大约 8%~12%磁盘写 IO 增加约 15%。主要开销来自 WAL 逻辑解码本身以及复制槽需要保留未消费的 WAL 段文件。如果你的源库本来负载就高我建议在低峰期开启管道并且设置max_slot_wal_keep_size来控制复制槽占用的磁盘空间上限。还有一点非常重要PG 的复制槽在没有消费者时WAL 文件不会被清理磁盘会持续增长直到撑满。所以管道意外停止后要尽快排查恢复而不是放着不管否则很容易出现迁移没做成、磁盘先爆了的事故。我在测试环境的默认配置下专门做了停掉 Worker 的测试12 小时未恢复时WAL 目录增长大约 3.8GB这是比较可观的数量级。5. 和 Flink CDC、Debezium 方案横向对比谁更适合湖仓实时场景用了 pg_duckpipe 一段时间后我得说句公道话它不是要全面替代 Flink CDC 或者 Debezium而是在轻量、快速、够用的场景下提供了第三条路。为了帮你做选型我把三条路线放在一张表格里对比一下。对比维度pg_duckpipeFlink CDC PipelineDebezium Kafka Iceberg Sink组件数量1 个 PG 扩展Flink 集群 CDC connector Sink connectorKafka Connect Debezium 自定义 Sink部署运维难度低SQL 创建扩展即可中高需要管理集群、作业、Checkpoint高组件多链路长吞吐能力单进程有限实测约 5000 rows/s 级别高分布式并行可扩展到十万级高Kafka 缓冲能力强延迟表现2~5 秒秒级到分钟级取决于 Flink 窗口和 checkpoint秒级到秒级以下取决于消费端类型转换支持依赖 DuckDB 映射新增类型要手动处理比较成熟的类型映射表Debezium 有丰富类型转换Schema 演进部分支持建议手动重建场景完善支持比较多Debezium 支持完整 Schema 演进全量增量一体化支持建管道时自带全量快照支持一般需要额外工具配合适用场景中小团队、单库轻量湖仓同步大规模实时数仓、复杂 ETL已有 Kafka 技术栈的企业一点个人理解。Debezium 最大的价值在于它是领域的事实标准几乎所有 Sink 都支持它Flink 适合那种我要在同步过程中做非常复杂的处理的场景比如多流 join、窗口聚合、数据质量校验而 pg_duckpipe 最适合的就是我不想处理任何额外集群只想要简单的表到表同步的人。我自己做选择时的一个判断标准是如果同步任务大于 20 条或者存在复杂的流处理需求我不会用 pg_duckpipe但如果只是把核心业务表增量搬到湖里做 ODSSQL 一把梭真的非常舒服。6. 排错和避坑经验我踩过的几个坑与解决办法最后这部分是纯实战经验都是文档里写得不清晰或者没写的地方。按重要程度排个序。6.1 复制槽冲突多个工具抢同一个 Slot 必然报错最容易踩的坑。我之前在一张测试表上同时跑了 Debezium 和 pg_duckpipe结果两个进程各自创建了自己的复制槽看起来互不干扰但实际上因为两个槽消费位点不同其中一个会长期保留大量旧 WAL 文件导致上游的 WAL 文件堆积。如果你确定只用 pg_duckpipe记得定期清理其他工具遗留的 slot-- 查看复制槽状态 SELECT slot_name, database, active, restart_lsn, confirmed_flush_lsn FROM pg_replication_slots; -- 确定不再使用的槽果断删除 SELECT pg_drop_replication_slot(旧槽名称);6.2 类型映射问题JSONB 和 NUMERIC 的意外表现我在建表时特意加了tags JSONB和amount NUMERIC(12,2)两个字段。同步完成后查询 Iceberg 表发现 JSONB 变成了 VARCHAR 长文本NUMERIC(12,2)变成了 DECIMAL(18, 4)。这个转换逻辑本身没错但如果下游写着amount 100.00这样的精确条件精度变化会导致匹配失败。解决方式是在建管道时指定一个自定义类型映射。pg_duckpipe 支持在 create_pipe 时传入column_types映射比如SELECT pg_duckpipe.create_pipe( pipe_name : user_orders_to_lake, source_table : public.user_orders, target_catalog : my_catalog, target_table : ods_user_orders, column_types : {amount: DECIMAL(12,2), tags: JSON} );这提醒我不要盲目相信默认类型映射尤其是有金融计算字段的表建管道以前最好先验证目标表结构和源表的一致性。6.3 Schema 变更加一列之后管道还在跑但新列写不进去业务上给表加了一列remark VARCHARIceberg 侧会自动跟着加上吗答案是否定的。pg_duckpipe 不会自动捕获 DDL 变更它只捕获 DML。加列后管道仍然会正常跑但新列的数据不会自动写入目标表。我见过很多人在这一步误以为工具坏了。实际上工具没坏只是你需要手动重建管道。如果你不想丢失已同步的历史数据推荐的操作是先暂停管道然后执行 ALTER TABLE 修改 Iceberg 表再恢复管道。顺序必须是这样否则先跑增量再 ALTER 会导致部分新列数据为空。6.4 大事务导致内存暴涨从 MyBatis 批量更新想到的前面提到批量更新 20 万行这其实是我踩过的一个真实的坑。有次我把一个在生产库跑批产生的 100 万行级 UPDATE 放进了测试环境结果 Worker 所在进程的内存占用直接飙到 1.8GB接近默认配置的上限。原因在于 pg_duckpipe 对一批变更事件的处理是先攒内存、再统一写文件大批量事务产生的所有变更都集中在同一批内存峰值就会特别高。如果你有大事务场景可以通过调低批量参数缓解或者限制源表的更新粒度。我在测试环境里把pg_duckpipe.batch_max_events调到 5000内存峰值降了约 40%但延迟略有上升。参数合适与否取决于你的具体数据模型没有绝对的标准答案。6.5 文件数目膨胀与 Compaction 策略冰berg 表最常见的健康度指标是平均文件大小和文件总数。pg_duckpipe 的合并机制不是无底的如果你的源表写入频率很高但每次量很小比如每秒几十条Iceberg 目录里还是会产生大量小文件查询性能会变差。我建议每运行一段时间后检查一次 Iceberg 表元数据-- 通过 duckdb 查看表的文件数和大小分布 SELECT * FROM duckdb.query($$ SELECT file_path, file_size_in_bytes FROM my_catalog.ods_user_orders.files ORDER BY file_size_in_bytes ASC LIMIT 20 $$);如果文件数很多且大小只有几 KB就需要跑一次 compaction。可以在my_catalog上用 DuckDB 的 Iceberg 扩展执行OPTIMIZE语句或者用 Spark 的 rewrite_data_files。Flink 有标准的 compaction job 可以挂着pg_duckpipe 则建议在低峰期手动跑避免自动压缩导致频繁的文件重写。6.6 停机恢复顺序千万别先删 slot假设某天服务器重启PG 恢复后会发现管道状态是 FAILED。这时候网上很多教程会让你先删掉旧 slot 再重建这是非常危险的做法——一旦 slot 删除未消费的 WAL 就被标记可清理而 Iceberg 表数据并没有补上就会造成数据空洞。正确的恢复顺序是检查pg_replication_slots中插件创建的 slot 是否还在。通过pg_duckpipe.pipes查看管道状态和最后 LSN。如果 slot 还在直接执行SELECT pg_duckpipe.resume_pipe(user_orders_to_lake)它会从保存的位点继续拉取。如果 slot 已经因为异常被清理了只能接受增量从当前时刻开始需要手动补一份全量快照到 Iceberg否则历史数据会缺一段。这个流程我吃过亏希望看到这篇文章的人走一遍测试环境的故障演练不然生产出问题的时候现场处理压力非常大。最后说两句掏心窝的话pg_duckpipe 目前还不适合和 Flink 在复杂流处理上一较高下它给我的感觉更像是数据工程师的瑞士军刀——不追求全场景覆盖但在PG 到 Iceberg中等数据量实时性要求不高这个区间里体验确实独一档。尤其是对团队里没有专职实时计算平台工程师的情况一个 DBA 或者后端开发就能通过几条 SQL 把管道搭起来后续运维成本几乎为零。我自己的建议是如果你的业务正好符合这个工具的适用范围不要犹豫先在一个低峰期的小表上测起来重点看三件事——你是否接受默认类型映射、源库负载是否吃得消逻辑解码开销、大事务场景下内存是否顶得住。测清楚这三件事再决定是否把它放到生产环境。至少在我这里它已经把一套原本需要三四个组件协作的同步链路压缩到了一次CREATE EXTENSION和一条建管道 SQL。
网站建设高端定制企业官网