Greenplum 到 Snowflake 数仓迁移

FreeGuideOnline 最新 2026-07-10

bash pg_dump --schema-only -t 'schema.table' -U user -h gp_host dbname > table_ddl.sql

根据映射表改写 DDL,关键调整点:
- **去掉分布键**:Snowflake 自身负责数据聚类,不需要 `DISTRIBUTED BY`
- **将分区表转换为 Snowflake 原生表**:Snowflake 使用微分区和自动聚类,无需显式 DDL 分区
- **修改表选项**:例如 `WITH (appendonly=true)` 等直接删除
- **索引移除**:Snowflake 自动优化,无需创建索引

示例重写:
```sql
-- Greenplum DDL
CREATE TABLE inventory (
    id bigint,
    product_name text,
    quantity int
) DISTRIBUTED BY (id)
PARTITION BY RANGE (quantity) (START (0) END (1000) EVERY (100));

转换为:

-- Snowflake DDL
CREATE OR REPLACE TABLE inventory (
    id BIGINT,
    product_name STRING,
    quantity INT
);

4.2 数据导出

推荐使用 COPY 将表数据导出为 CSV/Parquet 文件到对象存储或本地:

psql -h gp_host -U user -d dbname -c "\COPY (SELECT * FROM schema.table) TO '/data/table.csv' CSV HEADER"

大规模数据导出可并行执行多个 Segment 直连导出,或使用 GPLOAD 工具。

4.3 上传至云存储

  • AWS:使用 aws s3 cps3 sync 上传至 S3
  • Azureazcopy 到 Blob Storage
  • GCPgsutil cp 到 Cloud Storage

4.4 加载到 Snowflake

  1. 创建外部 Stage:
CREATE OR REPLACE STAGE my_stage
URL='s3://my-bucket/gp-data/'
CREDENTIALS=(AWS_KEY_ID='...' AWS_SECRET_KEY='...');
  1. 使用 COPY INTO 加载 CSV 文件:
COPY INTO inventory
FROM @my_stage/table.csv
FILE_FORMAT = (TYPE = CSV SKIP_HEADER = 1 FIELD_OPTIONALLY_ENCLOSED_BY='"')
ON_ERROR = 'CONTINUE';

更高效的方式:先将 CSV 转换为 Parquet,用 Parquet 加载获得更好的性能和压缩比。

4.5 数据验证

  • 行数对比:SELECT COUNT(*) FROM inventory
  • 校验和对比(对敏感列做 MD5 求和)
  • 业务规则抽样验证

5. 迁移存储过程与函数

5.1 处理 SQL 函数

  • Greenplum 中大量使用 PL/pgSQL、Python、R 等语言函数。
  • Snowflake 支持 JavaScript UDFSQL UDF 以及 Snowpark (Python/Java/Scala)
  • 遵循“先简后繁”原则,优先用 SQL UDF 重写简单逻辑。
  • 对于复杂过程,可整体重写为 Snowflake 存储过程(JavaScript 或 Snowpark)。

PL/pgSQL 示例重写

-- Greenplum
CREATE FUNCTION add_weekend_flag(d DATE) RETURNS BOOLEAN AS $$
BEGIN
    RETURN EXTRACT(DOW FROM d) IN (0,6);
END;
$$ LANGUAGE plpgsql;

迁移为 Snowflake SQL UDF:

CREATE OR REPLACE FUNCTION add_weekend_flag(d DATE)
RETURNS BOOLEAN
AS
$$
    DAYOFWEEKISO(d) IN (6,7)
$$;

5.2 外部查询与 FDW

Greenplum 常用外部表对接 HDFS、S3 等。Snowflake 可直接使用 External Table 或 Stage 直接查询,可实现类似功能。

6. 查询与脚本转换

6.1 常用函数差异速查

Greenplum Snowflake 等价
date_trunc('month', col) DATE_TRUNC('month', col)
age(timestamp1, timestamp2) DATEDIFF(day, timestamp2, timestamp1) 或自定义
generate_series(1,10) 使用 ROW_NUMBER() + TABLE(GENERATOR(...))
current_setting('TIMEZONE') CURRENT_TIMEZONE()
string_agg(col, ',') LISTAGG(col, ',')
regexp_split_to_table STRTOK_SPLIT_TO_TABLE
窗口函数 (ROW_NUMBER, RANK 等) 基本兼容,注意 Snowflake 需要显式 ORDER BY

6.2 优化查询模式

  • 去除 LIMIT + 大偏移量:使用 ROW_NUMBER() 配合 CTE 翻页,或利用 Snowflake 的查询结果缓存。
  • 利用 Snowflake 的 QUALIFY 简化窗口函数过滤:
SELECT * FROM sales
QUALIFY ROW_NUMBER() OVER (PARTITION BY sale_id ORDER BY created_at DESC) = 1;
  • 重写 DISTINCT ON:Greenplum 的 DISTINCT ON 在 Snowflake 中必须用 QUALIFY

7. 安全与权限迁移

7.1 用户与角色

  • Greenplum 角色模型与 PostgreSQL 类似。
  • Snowflake 权限层级:账户 → 数据库 → Schema → 表/视图。
  • 建议使用 RBAC(Role-Based Access Control)重新设计:
CREATE ROLE analyst;
GRANT USAGE ON DATABASE prod_db TO ROLE analyst;
GRANT USAGE ON SCHEMA public TO ROLE analyst;
GRANT SELECT ON ALL TABLES IN SCHEMA public TO ROLE analyst;

批量从 Greenplum 导出的角色需转化为 Snowflake 的 GRANT 语句。

7.2 行级安全与动态数据脱敏

  • Greenplum 中可能使用视图或 RLS 实现。
  • Snowflake 支持更先进的 Dynamic Data MaskingRow Access Policy,可做到列级与行列级安全。

8. 调度与 ETL 迁移

  • 将原有的基于 cron 或 Airflow 的任务中,Greenplum 的 SQL 直接替换为 Snowflake SQL。
  • 使用 Snowflake 的 TASKSTREAMS 实现 ETL 内部调度,与外部 Airflow 配合。
  • Snowpipe 可实现文件到达自动摄入,减少调度负担。

9. 性能验证与优化

9.1 选择合适的虚拟仓库

  • 数据加载使用 X-Small 至 Medium 的仓库,监测加载吞吐。
  • 查询按复杂度设置不同大小的仓库,利用多集群并发处理高并发。

9.2 监控查询性能

使用 QUERY_HISTORYINFORMATION_SCHEMA 分析过去查询:

SELECT query_text, total_elapsed_time
FROM SNOWFLAKE.ACCOUNT_USAGE.QUERY_HISTORY
WHERE execution_status = 'SUCCESS'
ORDER BY total_elapsed_time DESC;

9.3 启用自动聚类与物化视图

  • 对大表设置自动聚类键以提升查询效率:
ALTER TABLE large_table CLUSTER BY (date, region);