探索 Snowflake 与 Postgres 之间的双向数据流动模式 | 技术实践
日程已上线 100%!6.26-27 AICon 上海站: 13大专题+1个动手实验室、近60场重磅议题,集结清华、复旦等知名高校教授及阿里、腾讯、字节、小红书、Google Cloud等头部企业技术专家,围绕Agent工程化落地等相关议题展开分享, 点击了解详情
_2026 年,智能体将在企业级应用中取得哪些实质性突破?_ 点击下载 _《2026 年 AI 与数据发展预测》白皮书,获悉专家一手前瞻,抢先拥抱新的工作方式!_
现代数据架构越来越要求事务型与分析型工作负载能够无缝共存。PostgreSQL 仍然是事务型应用的核心基础设施——支撑电商订单处理、实时库存系统以及面向客户的 API——而 Snowflake 则是所有分析型数据与 AI 的基础平台。挑战在于:如何让数据在这两个系统之间双向可靠流动,同时将延迟与运维开销降到最低。
历史上,要打通 OLTP 与 OLAP 系统,需要拼接外部 ETL 工具、管理云存储桶、配置 IAM 角色,并维护脆弱的 CDC 数据管道。团队花在基础设施“打通”上的时间,往往超过从数据中产生价值的时间。
本文的目标是探索如何利用 Snowflake 的一些最新创新能力,在 Postgres 与 Snowflake 之间支持五种关键的数据流动模式。如下所示:
有哪些变化?
近期 Snowflake 的一些产品能力显著简化了这一问题:
- Snowflake Postgres(PuPr) :一种完全托管的 PostgreSQL 服务,原生运行在 Snowflake 生态系统中。它消除了外部 PostgreSQL 托管的需求,同时与 Snowflake 数据平台实现一流集成;
- pg_lake(PuPr) :PostgreSQL 的扩展能力,允许在 Postgres 内直接创建 Apache Iceberg 表。在 Postgres 中写入的数据,可以通过共享 Iceberg 元数据被 Snowflake 直接查询——无需文件导出、无需中转存储、无需 ETL 管道;
- pg_incremental(PuPr) :用于调度式增量同步的扩展组件。与 pg_lake 结合使用时,可以实现轻量级 CDC,只同步发生变化的数据行;
- Snowflake 托管 Iceberg 存储(PuPr) :Postgres 管理的 Iceberg 表使用 Snowflake 内部存储,并通过托管凭证访问。无需外部 S3、无需 IAM、无需存储集成配置;
- Openflow(正式发布) :Snowflake 的托管数据集成平台(基于 Apache NiFi 构建),提供预置的 CDC 连接能力,包括基于 PostgreSQL WAL 的变更捕获,以及通过 Snowpipe Streaming 进行数据传输。
这些能力共同构建了一个完整的数据流动能力体系——从简单的批量加载,到实时 CDC——全部在一个统一平台内完成。
本文中的所有模式都使用 ORDERS 表作为示例数据源。该表模拟典型电商订单生命周期,约 28,000 行数据。本例中数据来源于 SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.Orders。
CREATE TABLE cdc_demo.orders ( order_id BIGINT PRIMARY KEY, customer_id BIGINT, order_status VARCHAR(1), total_price DECIMAL(15,2), order_date DATE, order_priority VARCHAR(15), clerk VARCHAR(15), ship_priority INTEGER, comment VARCHAR(79), created_at TIMESTAMP DEFAULT NOW(), updated_at TIMESTAMP DEFAULT NOW());
复制代码
模式 1:批量数据流动 —— Postgres 到 Snowflake
业务场景: 某零售公司在 PostgreSQL 电商系统中全天处理订单。每天夜间,分析团队需要在 Snowflake 中获取所有订单的完整快照,用于报表分析、需求预测以及财务对账;
技术方案: 在 Postgres 中创建 Iceberg 表。“USING Iceberg”子句允许 pg_lake 在 Postgres 内创建并写入 Iceberg 表。然后通过 INSERT/SELECT 将订单数据写入该表。
在 Snowflake 侧,创建目录集成(Catalog Integration)后再创建 Iceberg 表。这些操作属于元数据层操作,不涉及实际数据搬运。当 Iceberg 表创建完成后,即可在 Snowflake 中直接查询。
步骤 1:在 Postgres 中启用 pg_lake
-- Connect to Snowflake Postgres instanceCREATE EXTENSION IF NOT EXISTS pg_lake CASCADE;
复制代码
步骤 2:创建 Iceberg 表并批量加载数据
-- Create an Iceberg table and bulk load all orders into itCREATE TABLE cdc_demo.orders_iceberg ( order_id BIGINT, customer_id BIGINT, order_status VARCHAR(1), total_price DECIMAL(15,2), order_date DATE, order_priority VARCHAR(15), clerk VARCHAR(15), ship_priority INTEGER, comment VARCHAR(79), created_at TIMESTAMP, updated_at TIMESTAMP) USING iceberg;
-- Bulk load from the source tableINSERT INTO cdc_demo.orders_icebergSELECT * FROM cdc_demo.orders;
复制代码
步骤 3:在 Snowflake 中创建目录集成
-- In Snowflake: create a catalog integration pointing to the Postgres instanceCREATE OR REPLACE CATALOG INTEGRATION pg_orders_catalog CATALOG_SOURCE = SNOWFLAKE_POSTGRES TABLE_FORMAT = ICEBERG CATALOG_NAMESPACE = 'cdc_demo' REST_CONFIG = ( POSTGRES_INSTANCE = 'Snowflake_Postgres_Demo' CATALOG_NAME = 'postgres' ACCESS_DELEGATION_MODE = VENDED_CREDENTIALS ) ENABLED = TRUE;
复制代码
步骤 4:在 Snowflake 创建 Iceberg 表
-- Create the Snowflake Iceberg table referencing the Postgres-managed Iceberg dataCREATE OR REPLACE ICEBERG TABLE orders_iceberg CATALOG = 'pg_orders_catalog' CATALOG_TABLE_NAME = 'orders_iceberg' CATALOG_NAMESPACE = 'cdc_demo' AUTO_REFRESH = TRUE;
复制代码
步骤 5:查询验证数据
SELECT COUNT(*) FROM orders_iceberg;
SELECT order_status, COUNT(*), SUM(total_price) AS total_revenueFROM orders_icebergGROUP BY order_status;
复制代码
模式 2:批量数据流动 —— Snowflake 到 Postgres
业务场景: 数据科学团队在 Snowflake 中构建订单优先级预测模型,结果需要回写到 Postgres,使业务系统能够实时展示预测结果;
技术方案: 在 Snowflake 中生成结果表后,将数据写入 stage(Parquet 文件)。Postgres 从 stage 拉取数据并写入本地表。
步骤 1:创建存储集成
-- In Snowflake: create a storage integration for the Postgres managed storageCREATE OR REPLACE STORAGE INTEGRATION pg_stage_integration TYPE = POSTGRES_INTERNAL_STORAGE POSTGRES_INSTANCE = 'Snowflake_Postgres_Demo';
-- Create a stage using this integrationCREATE OR REPLACE STAGE pg_orders_stage RELATIVE_URL = '/orders_export' STORAGE_INTEGRATION = pg_stage_integration;
复制代码
步骤 2:导出数据到 stage
COPY INTO @pg_orders_stage/orders_bulk_FROM ( SELECT order_id, customer_id, order_status, total_price, order_date, order_priority, clerk, ship_priority, comment FROM orders_iceberg)FILE_FORMAT = (TYPE = PARQUET)HEADER = TRUEOVERWRITE = TRUE;
复制代码
步骤 3:Postgres 读取数据
-- On Postgres: create the destination tableCREATE TABLE cdc_demo.orders_from_snowflake ( order_id BIGINT PRIMARY KEY, customer_id BIGINT, order_status VARCHAR(1), total_price DECIMAL(15,2), order_date DATE, order_priority VARCHAR(15), clerk VARCHAR(15), ship_priority INTEGER, comment VARCHAR(79));
-- Load the Parquet files from the stage into the Postgres tableCOPY cdc_demo.orders_from_snowflakeFROM '@STAGE/orders_export/orders_bulk_*.parquet';
复制代码
步骤 4:查询验证
SELECT COUNT(*) FROM cdc_demo.orders_from_snowflake;-- Returns: 28,373
SELECT order_status, COUNT(*)FROM cdc_demo.orders_from_snowflakeGROUP BY order_status;
复制代码
模式 3:CDC —— 从 Postgres 到 Snowflake
业务场景: 一个电商应用持续处理新订单、更新订单状态(已发货、已送达、已退货)以及取消订单。分析团队需要在几分钟内将这些变化同步到 Snowflake,而不是等待数小时,以支持展示订单履约指标与营收追踪的实时仪表盘;
技术方案: 这代表典型的 OLTP → OLAP 模式,即在源数据库中检测数据变更,并将其传递到分析系统。在该模式中,Postgres 源表上的插入与更新操作会被捕获。系统每分钟检测一次变更,并将其写入 Iceberg 表。在 Snowflake 侧,会创建 catalog integration 和 Iceberg 表,并通过每分钟一次的 REFRESH 进行轮询,从而使最新变更可以被查询访问。
步骤 1:在 Postgres 上启用 pg_incremental 和 pg_cron
CREATE EXTENSION IF NOT EXISTS pg_cron;CREATE EXTENSION IF NOT EXISTS pg_incremental CASCADE;
复制代码
步骤 2:创建 Iceberg 目标表并执行初始批量加载
-- Create an Iceberg table for CDC dataCREATE TABLE cdc_demo.orders_iceberg_cdc ( order_id BIGINT, customer_id BIGINT, order_status VARCHAR(1), total_price DECIMAL(15,2), order_date DATE, order_priority VARCHAR(15), clerk VARCHAR(15), ship_priority INTEGER, comment VARCHAR(79), updated_at TIMESTAMP) USING iceberg;
-- Initial bulk load of current dataINSERT INTO cdc_demo.orders_iceberg_cdcSELECT order_id, customer_id, order_status, total_price, order_date, order_priority, clerk, ship_priority, comment, updated_atFROM cdc_demo.orders;
复制代码
步骤 3:创建时间间隔 pipeline,用于检测并同步变更