返回
RSS Snowflake Engineering (Medium) AI 逐段翻译 精选 发布 2026-08-31 22:01 收录于 09-01

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 节
  1. 什么是 Snowflake Postgres?
  2. 主要能力:
  3. 步骤 1:创建 Postgres 实例
  4. 网络策略设置
  5. 步骤 2:构建应用 schema
  6. 步骤 3:数据镜像——持续 CDC 复制
  7. 镜像工作原理
  8. 创建镜像
  9. 创建的内容
  10. 步骤 4:查询镜像数据和 CDC
  11. 基表(当前状态)
  12. $CHANGES 表(CDC 审计日志)
  13. $LIVE 视图(实时)
  14. schema 演进
  15. 步骤 5:pg_lake——选择性导出为 Apache Iceberg
  16. 镜像与 pg_lake 对比
  17. 模式 1:共享 Iceberg(无需设置 S3)
  18. 模式 2:客户管理的 S3(完全控制)
  19. 步骤 6:目录链接数据库(自动发现)
  20. 步骤 7:成本管理
  21. 成本监控
  22. 何时使用什么
  23. 关键要点
  24. 下一步是什么?

译文

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 上,人们通过突出显示和回应这个故事来继续对话。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存