← 文章 / 科技资讯
infoq 2026/6/24 · 2026-06-24 18:00:38 · 0 阅读

探索 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,用于检测并同步变更

原始来源: infoq

评论 (0)