探索 Snowflake 与 Postgres 之间的双向数据流动模式 | 技术实践

InfoQ 中文 2026-06-23T09:13:53.654206

探索 Snowflake 与 Postgres 之间的双向数据流动模式 | 技术实践

现代数据架构要求事务型与分析型工作负载无缝共存,PostgreSQL 处理交易,Snowflake 承载分析。本文介绍如何利用 Snowflake 最新能力在两者之间实现双向数据流动,涵盖批量同步与实时 CDC 等关键模式。

现代数据架构越来越要求事务型与分析型工作负载能够无缝共存。PostgreSQL 依然是事务型应用的核心基础设施——支撑电商订单处理、实时库存系统以及面向客户的 API——而 Snowflake 则是所有分析型数据与 AI 的基础平台。核心挑战在于:如何让数据在这两个系统之间双向可靠流动,同时将延迟与运维开销降至最低。

以往,打通 OLTP 与 OLAP 系统需要拼接外部 ETL 工具、管理云存储桶、配置 IAM 角色,并维护脆弱的 CDC 数据管道。团队花在基础设施“打通”上的时间,往往超过从数据中提取价值的时间。

本文将探索如何利用 Snowflake 的最新创新能力,在 Postgres 与 Snowflake 之间支持五种关键的数据流动模式。如下图所示:

Snowflake 的关键变化

近期 Snowflake 推出的一系列产品能力显著简化了这一问题:

这些能力共同构建了一个完整的数据流动体系——从简单的批量加载到实时 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 的增量同步功能。)

关键要点

查看原文