探索 Snowflake 与 Postgres 之间的双向数据流动模式 | 技术实践
探索 Snowflake 与 Postgres 之间的双向数据流动模式 | 技术实践
现代数据架构要求事务型与分析型工作负载无缝共存,PostgreSQL 处理交易,Snowflake 承载分析。本文介绍如何利用 Snowflake 最新能力在两者之间实现双向数据流动,涵盖批量同步与实时 CDC 等关键模式。
现代数据架构越来越要求事务型与分析型工作负载能够无缝共存。PostgreSQL 依然是事务型应用的核心基础设施——支撑电商订单处理、实时库存系统以及面向客户的 API——而 Snowflake 则是所有分析型数据与 AI 的基础平台。核心挑战在于:如何让数据在这两个系统之间双向可靠流动,同时将延迟与运维开销降至最低。
以往,打通 OLTP 与 OLAP 系统需要拼接外部 ETL 工具、管理云存储桶、配置 IAM 角色,并维护脆弱的 CDC 数据管道。团队花在基础设施“打通”上的时间,往往超过从数据中提取价值的时间。
本文将探索如何利用 Snowflake 的最新创新能力,在 Postgres 与 Snowflake 之间支持五种关键的数据流动模式。如下图所示:

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()
);
模式一:批量数据流动 —— Postgres 到 Snowflake
业务场景:某零售公司在 PostgreSQL 电商系统中全天处理订单。每天夜间,分析团队需要在 Snowflake 中获取所有订单的完整快照,用于报表分析、需求预测以及财务对账。
技术方案:在 Postgres 中使用 USING Iceberg 子句创建 Iceberg 表,pg_lake 负责在 Postgres 内创建并写入 Iceberg 表。然后通过 INSERT/SELECT 将订单数据写入该表。
在 Snowflake 侧,先创建目录集成(Catalog Integration),再创建 Iceberg 表。这些操作属元数据层操作,不涉及实际数据搬运。Iceberg 表创建完成后,即可直接在 Snowflake 中查询。
步骤 1:在 Postgres 中启用 pg_lake
-- Connect to Snowflake Postgres instance
CREATE EXTENSION IF NOT EXISTS pg_lake CASCADE;
步骤 2:创建 Iceberg 表并批量加载数据
-- Create an Iceberg table and bulk load all orders into it
CREATE 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 table
INSERT INTO cdc_demo.orders_iceberg
SELECT * FROM cdc_demo.orders;
步骤 3:在 Snowflake 中创建目录集成
-- In Snowflake: create a catalog integration pointing to the Postgres instance
CREATE 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 data
CREATE 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_revenue
FROM orders_iceberg
GROUP BY order_status;
模式二:批量数据流动 —— Snowflake 到 Postgres
业务场景:数据科学团队在 Snowflake 中构建订单优先级预测模型,结果需要回写到 Postgres,使业务系统能够实时展示预测结果。
技术方案:在 Snowflake 中生成结果表后,将数据写入 stage(Parquet 文件)。Postgres 从 stage 拉取数据并写入本地表。

步骤 1:创建存储集成
-- In Snowflake: create a storage integration for the Postgres managed storage
CREATE OR REPLACE STORAGE INTEGRATION pg_stage_integration
TYPE = POSTGRES_INTERNAL_STORAGE
POSTGRES_INSTANCE = 'Snowflake_Postgres_Demo';
-- Create a stage using this integration
CREATE 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 = TRUE
OVERWRITE = TRUE;
步骤 3:Postgres 读取数据
-- On Postgres: create the destination table
CREATE 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 table
COPY cdc_demo.orders_from_snowflake
FROM '@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_snowflake
GROUP BY order_status;
模式三:CDC —— 从 Postgres 到 Snowflake
业务场景:一个电商应用持续处理新订单、更新订单状态(已发货、已送达、已退货)以及取消订单。分析团队需要在几分钟内将这些变化同步到 Snowflake,而不是等待数小时,以支持展示订单履约指标与营收追踪的实时仪表盘。

技术方案:这代表典型的 OLTP → OLAP 模式,即在源数据库中检测数据变更,并将其传递到分析系统。此处,Postgres 源表上的插入与更新操作会被捕获。系统每分钟检测一次变更,并将其写入 Iceberg 表。在 Snowflake 侧,创建目录集成和 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 data
CREATE 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),
created_at TIMESTAMP,
updated_at TIMESTAMP
) USING iceberg;
-- Initial bulk load (后续增量同步将通过 pg_incremental 实现)
INSERT INTO cdc_demo.orders_iceberg_cdc
SELECT * FROM cdc_demo.orders;
(注:后续增量配置步骤因原文截断未完整列出,实际应用需配合 pg_cron 定时任务和 pg_incremental 的增量同步功能。)
关键要点
- 统一平台简化集成:Snowflake Postgres(PuPr)与 pg_lake、pg_incremental 等扩展将 Postgres 和 Snowflake 的数据流动完全纳入一个管理生态,消除了外部 ETL 工具和复杂配置。
- Iceberg 表实现无感共享:通过 Iceberg 格式,Postgres 写入的数据可直接被 Snowflake 查询,无需导出和复制;反之,Snowflake 的输出也能通过托管 Stage 轻松回传至 Postgres。
- 增量 CDC 降低延迟:pg_incremental 结合 pg_lake 可实现分钟级的变更数据捕获,满足实时仪表盘等低延迟业务需求,且无需部署独立的 CDC 基础设施。
- 批量与实时灵活组合:无论夜间批量同步还是持续增量同步,均可通过统一的产品能力按需选择,降低运维成本。
- 安全与托管能力:Snowflake 托管 Iceberg 存储和 Vended Credentials 机制消除了外部存储集成和安全配置的复杂性,数据流动更加安全可控。