Snowflake Postgres:从OLTP到分析的无ETL实战指南
DataHot 速览
Snowflake Postgres 是运行在 Snowflake 账户内的全托管 PostgreSQL 18 服务,支持 Data Mirroring 实现亚分钟级 CDC 持续复制,并可通过 pg_lake 将数据导出为 Apache Iceberg 表到客户自有 S3。文章用电商场景的端到端 POC,演示了创建实例、配置网络策略、构建 schema、启用连续复制、查询 $CHANGES/$LIVE 对象等完整流程。所有代码均来自真实实现,强调无需管理基础设施、无需外部 CDC 连接器,且数据格式开放、无厂商锁定。
为什么值得关注:数据从业者可借此了解 Snowflake 将 OLTP 与分析打通的 Zero-ETL 路径,以及用 Apache Iceberg 保持开放格式的具体做法,对评估托管 PG 与湖仓集成有直接参考价值。
本文目录 24 节
- 什么是 Snowflake Postgres?
- 主要能力:
- 步骤 1:创建 Postgres 实例
- 网络策略设置
- 步骤 2:构建应用 schema
- 步骤 3:数据镜像——持续 CDC 复制
- 镜像工作原理
- 创建镜像
- 创建的内容
- 步骤 4:查询镜像数据和 CDC
- 基表(当前状态)
- $CHANGES 表(CDC 审计日志)
- $LIVE 视图(实时)
- schema 演进
- 步骤 5:pg_lake——选择性导出为 Apache Iceberg
- 镜像与 pg_lake 对比
- 模式 1:共享 Iceberg(无需设置 S3)
- 模式 2:客户管理的 S3(完全控制)
- 步骤 6:目录链接数据库(自动发现)
- 步骤 7:成本管理
- 成本监控
- 何时使用什么
- 关键要点
- 下一步是什么?
译文
AI 逐段翻译一份关于 Snowflake Postgres、数据镜像、pg_lake 和 Apache Iceberg 的动手实践指南——以零基础设施将事务数据引入 Snowflake。— 作者:Dipal Mahajan | Snowflake 高级解决方案架构师
现代数据架构要求分析系统与操作数据库保持同步——实时同步,无需复杂的 ETL 管道、外部 CDC 连接器,也不受数据格式供应商锁定。Snowflake Postgres 正好提供了这一点。
Snowflake Postgres 是一项完全托管的 PostgreSQL 服务,在您的 Snowflake 账户内运行。它不仅仅是另一个托管 Postgres——它内置了数据移动功能,消除了 OLTP 和 OLAP 之间的传统差距。写入 Postgres,从 Snowflake 查询——亚分钟级延迟,零基础设施需要管理。
在本文中,我将演示一个端到端的概念验证:创建 Postgres 实例、构建电商模式、通过数据镜像启用连续 CDC 复制、探索 $CHANGES 和 $LIVE 对象,最后使用 pg_lake 将数据以 Apache Iceberg 表形式导出到客户管理的 S3。每个代码块都来自真实可运行的实现。

什么是 Snowflake Postgres?
Snowflake Postgres 是一项托管的 PostgreSQL 18 服务,直接在您的 Snowflake 账户内运行。它为您的 OLTP 工作负载提供完全的 PostgreSQL 兼容性,同时提供原生数据移动功能,将该数据引入 Snowflake 进行分析。
主要能力:
- 完全 PostgreSQL 兼容性 — 使用 psql、pgAdmin、DBeaver 或任何 PG 客户端
- 托管的计算和存储 — 无需修补、无需 vacuum 调优、无需管理复制槽
- 内置高可用性和自动故障转移
- 数据镜像 — 持续 CDC 复制到 Snowflake。
- pg_lake — 选择性数据导出为 Apache Iceberg 到 S3(开放格式,无供应商锁定)
- schema 演进 — ALTER TABLE 自动传播到 Snowflake
步骤 1:创建 Postgres 实例
我通过 Snowsight 使用以下配置创建了一个 Postgres 实例:

网络策略设置
Snowflake Postgres 要求在任何外部客户端连接之前有明确的网络策略。否则,默认情况下,端口 5432 上的所有连接都会被阻止。
-- Create network rule
CREATE OR REPLACE NETWORK RULE SECURITY_NETWORK_DB.PUBLIC.postgres_poc_ingress_rule
MODE = POSTGRES_INGRESS
TYPE = IPV4
VALUE_LIST = ('IP_Ranges');
-- Create and attach policy
CREATE OR REPLACE NETWORK POLICY postgres_poc_policy
ALLOWED_NETWORK_RULE_LIST = ('SECURITY_NETWORK_DB.PUBLIC.postgres_poc_ingress_rule');
ALTER POSTGRES INSTANCE "Postgres_POC" SET NETWORK_POLICY = 'POSTGRES_POC_POLICY';步骤 2:构建应用 schema
通过 psql 连接后,我创建了一个简单的电商 schema 来演示完整的数据移动生命周期:
— 从终端连接: psql “postgres://snowflake_admin:<PASSWORD>@<HOST>:5432/postgres?sslmode=require”

-- Create tables
CREATE TABLE customers (
customer_id SERIAL PRIMARY KEY,
first_name VARCHAR(50) NOT NULL,
last_name VARCHAR(50) NOT NULL,
email VARCHAR(100) UNIQUE NOT NULL,
phone VARCHAR(20),
created_at TIMESTAMP DEFAULT NOW()
);
CREATE TABLE products (
product_id SERIAL PRIMARY KEY,
name VARCHAR(200) NOT NULL,
price NUMERIC(10,2) NOT NULL,
stock_quantity INTEGER DEFAULT 0,
category VARCHAR(50)
);
CREATE TABLE orders (
order_id SERIAL PRIMARY KEY,
customer_id INTEGER REFERENCES customers(customer_id),
status VARCHAR(20) DEFAULT 'pending',
total_amount NUMERIC(10,2)
);
--Insert sample data:
INSERT INTO customers (first_name, last_name, email, phone) VALUES
('Rahul', 'Sharma', '[email protected]', '+91–9876543210'),
('Priya', 'Patel', '[email protected]', '+91–9876543211'),
('Amit', 'Kumar', '[email protected]', '+91–9876543212');
INSERT INTO products (name, price, stock_quantity, category) VALUES
('Wireless Headphones', 4999.00, 50, 'Electronics'),
('Cotton T-Shirt', 699.00, 200, 'Apparel'),
('Python Programming Book', 549.00, 75, 'Books');步骤 3:数据镜像——持续 CDC 复制
数据镜像是头条功能。它持续将每个 INSERT、UPDATE、DELETE 和 schema 更改从 Postgres 复制到 Snowflake——自动完成,无需连接器、无需 Kafka、无需 Debezium。只需一个 SQL 调用即可设置。
镜像工作原理
在底层,镜像使用 PostgreSQL 的逻辑复制(WAL 解码)结合 pg_lake 的开源 Iceberg 写入器。Postgres 上的 CDC worker 捕获更改并将其写入每个表的 $CHANGES Iceberg 表。然后,Snowflake 无服务器任务(APPLY)按可配置间隔将这些更改合并到目标表中。
创建镜像
-- Grant permissions
GRANT APPLICATION ROLE snowflake.postgres_mirror_admin TO ROLE "ACCOUNTADMIN";
GRANT USAGE ON POSTGRES INSTANCE "Postgres_POC" TO APPLICATION snowflake;
-- Create the mirror (30-second refresh interval - minimum)
CALL SNOWFLAKE.POSTGRES.CREATE_MIRROR(
mirror_name => 'snowflake_postgres_mirror',
postgres_instance => 'Postgres_POC',
postgres_database => 'postgres',
target_database => 'SNOWFLAKE_POSTGRES_MIRROR',
postgres_tables => ['public.customers', 'public.products',
'public.orders', 'public.order_items'],
postgres_schemas => NULL,
refresh_interval => '30 seconds'
);创建的内容
对于每个镜像表,Snowflake 在目标数据库中创建三个对象:
例如:Customers 表


步骤 4:查询镜像数据和 CDC
基表(当前状态)

$CHANGES 表(CDC 审计日志)
$CHANGES 表是一个仅追加的 Iceberg 表,记录来自 Postgres 的每次更改,并带有完整的事务元数据:

$LIVE 视图(实时)
$LIVE 视图将基表与 $CHANGES 中未合并的更改结合起来,让您获得最新的状态(约 30 秒延迟),而无需等待 APPLY 任务:

schema 演进
然后我通过在 Postgres 中添加一列来测试 schema 演进:
-- In Postgres:
ALTER TABLE customers ADD COLUMN city VARCHAR(100);
INSERT INTO customers (first_name, last_name, email, phone, city) VALUES
('Sneha', 'Joshi', '[email protected]', '+91-9876543216', 'Mumbai'),
('Ravi', 'Verma', '[email protected]', '+91-9876543217', 'Delhi');步骤 5:pg_lake——选择性导出为 Apache Iceberg
镜像提供整个表的持续、自动复制,而 pg_lake 提供选择性、按需的数据移动。您可以选择导出什么、何时导出、导出到哪里——以开放的 Apache Iceberg 格式。
镜像与 pg_lake 对比

模式 1:共享 Iceberg(无需设置 S3)
最简单的模式——Snowflake 自动管理存储。无需 S3 存储桶,无需 IAM 角色。
-- In Postgres: Enable pg_lake and create Iceberg table
CREATE EXTENSION IF NOT EXISTS pg_lake CASCADE;
CREATE TABLE customers_iceberg (
customer_id INTEGER,
first_name TEXT,
last_name TEXT,
email TEXT,
phone TEXT,
city TEXT,
created_at TIMESTAMP
) USING iceberg;
-- Export data
INSERT INTO customers_iceberg
SELECT customer_id, first_name, last_name, email, phone, city, created_at
FROM customers;
-- In Snowflake: Create catalog integration and read the data
CREATE OR REPLACE CATALOG INTEGRATION pg_lake_catalog
CATALOG_SOURCE = SNOWFLAKE_POSTGRES
TABLE_FORMAT = ICEBERG
CATALOG_NAMESPACE = 'public'
REST_CONFIG = (
POSTGRES_INSTANCE = 'Postgres_POC'
CATALOG_NAME = 'postgres'
ACCESS_DELEGATION_MODE = VENDED_CREDENTIALS
)
ENABLED = TRUE;
-- Query directly from Snowflake!
CREATE OR REPLACE ICEBERG TABLE DEV.PUBLIC.CUSTOMERS_ICEBERG
CATALOG = 'PG_LAKE_CATALOG'
CATALOG_TABLE_NAME = 'customers_iceberg'
AUTO_REFRESH = TRUE;
SELECT * FROM DEV.PUBLIC.CUSTOMERS_ICEBERG ORDER BY CUSTOMER_ID;模式 2:客户管理的 S3(完全控制)
当您需要将数据放在您的 S3 存储桶上——可由 Spark、Trino 或任何与 Iceberg 兼容的工具读取——请使用客户管理的存储模式:
-- In Snowflake: Create storage integration
CREATE STORAGE INTEGRATION pg_lake_storage_int
TYPE = POSTGRES_EXTERNAL_STORAGE
STORAGE_PROVIDER = 'S3'
STORAGE_AWS_ROLE_ARN = 'arn:aws:iam::<your-account-id>:role/pg_lake_role'
ENABLED = TRUE
STORAGE_ALLOWED_LOCATIONS = ('s3://<your-bucket>/pg-lake/');
-- Attach to instance
ALTER POSTGRES INSTANCE "Postgres_POC"
SET STORAGE_INTEGRATION = 'PG_LAKE_STORAGE_INT';
-- In Postgres: Create Iceberg table on YOUR S3
CREATE FOREIGN TABLE customers_s3 (
customer_id INTEGER,
first_name TEXT,
last_name TEXT,
email TEXT,
phone TEXT,
city TEXT,
created_at TIMESTAMP
) SERVER pg_lake_iceberg
OPTIONS (location 's3://<your-bucket>/pg-lake/customers');
-- Export data to YOUR bucket
INSERT INTO customers_s3
SELECT customer_id, first_name, last_name, email, phone, city, created_at
FROM customers;
-- Verify (reading back from S3)
SELECT * FROM customers_s3 ORDER BY customer_id;步骤 6:目录链接数据库(自动发现)
无需在 Snowflake 中逐个创建 Iceberg 表,目录链接数据库会自动发现来自 Postgres 的所有 pg_lake Iceberg 表——并保持同步。
CREATE DATABASE PG_LAKE_DB
LINKED_CATALOG = (
CATALOG = PG_LAKE_CATALOG
ALLOWED_WRITE_OPERATIONS = NONE
);
-- Automatically discovers all Iceberg tables!
SHOW ICEBERG TABLES IN DATABASE PG_LAKE_DB;
-- Query without any manual table creation
SELECT * FROM PG_LAKE_DB.PUBLIC.CUSTOMERS_ICEBERG;
SELECT * FROM PG_LAKE_DB.PUBLIC.PRODUCTS_ICEBERG;

步骤 7:成本管理
镜像成本有两个主要组成部分:
- 计算用于将更改应用到目标表的无服务器任务,按标准 Snowflake 无服务器任务费率计费。
- 存储分为源 Postgres 实例(通过 pg_lake 保存 $changes 和 metalog Iceberg 表)和目标 Snowflake 数据库(保存物化目标表),按标准 Snowflake 存储费率计费。
成本监控
-- Monitor mirroring compute cost
SELECT * FROM SNOWFLAKE.ACCOUNT_USAGE.SERVERLESS_TASK_HISTORY
WHERE TASK_NAME LIKE 'APPLY_MIRROR_%' ORDER BY START_TIME DESC;
-- Monitor Postgres storage
SELECT * FROM SNOWFLAKE.ACCOUNT_USAGE.POSTGRES_STORAGE_USAGE_HISTORY
ORDER BY START_TIME DESC;何时使用什么
黄金法则:如果您需要自动更新和删除 → 使用镜像。如果您需要将数据放在您的 S3 上并以开放格式呈现 → 使用 pg_lake。您可以同时使用两者。
关键要点
零基础设施: 无需 Kafka、无需 Debezium、无需 Airbyte、无需外部 CDC 服务。镜像完全在 Snowflake 内部的无服务器计算上运行。
亚分钟级延迟: $LIVE 视图提供约 30 秒的新鲜度,无论合并间隔如何。您等待 Postgres 数据在 Snowflake 中可查询的时间永远不会超过 30 秒。
开放格式(Apache Iceberg): pg_lake 写入标准 Iceberg 表。您的数据以 Parquet 格式存储在 S3 上——可由 Spark、Trino、Flink 或任何工具读取。无供应商锁定。
schema 演进: 在 Postgres 中执行 ALTER TABLE ADD/DROP COLUMN 会自动即时传播到 Snowflake。无需重新配置。
完整的带审计追踪的 CDC:$CHANGES 提供 7 天可查询的变更流,并具有事务排序。每条 INSERT、UPDATE、DELETE 都会记录 Postgres 提交元数据。
成本高效:无服务器 APPLY 任务,无需常驻连接器进程。不使用时挂起 Postgres 实例。$CHANGES 在 7 天后自动清除。
下一步是什么?
- 动态表: 在镜像表之上构建增量管道,用于下游转换。
- Cortex AI: 使用 Snowflake 的 AI 函数对刚复制的 Postgres 数据进行情感分析、分类等。
- pg_incremental: 为仅追加工作负载(物联网、事件、日志)设置自动化的、计划性的 Iceberg 导出。
- Snowflake CoWork: 将 Snowflake 的对话式 AI 指向您的镜像数据库,进行自然语言分析。
更多信息: https://docs.snowflake.com/en/user-guide/snowflake-postgres/
Snowflake Postgres:从 OLTP 到分析的零 ETL 最初发布在 Snowflake Builders Blog: Data Engineers, App Developers, AI, & Data Science 在 Medium 上,人们通过突出显示和回应这个故事来继续对话。
这篇内容对你有用吗?
反馈只用于改善内容筛选,不等同于收藏