Snowflake 开发

snowflake-development
分类编程
作者Alireza Rezvani
许可MIT
评分4.20/5
使用10.9K

Snowflake 开发

涵盖 Snowflake SQL、数据管道、Cortex AI 和 Snowpark Python 开发。包括冒号前缀规则、半结构化数据、MERGE upserts、Dynamic Tables、Streams+Tasks、Cortex AI 函数、Agent 规范、性能调优和安全加固。

> 最初由 James Cha-Earley 贡献 —— 由 claude-skills 团队增强并集成。

快速上手

bash
# 生成 MERGE upsert 模板
python scripts/snowflake_query_helper.py merge --target customers --source staging_customers --key customer_id --columns name,email,updated_at

生成 Dynamic Table 模板

python scripts/snowflake_query_helper.py dynamic-table --name cleaned_events --warehouse transform_wh --lag "5 minutes"

生成 RBAC 授权语句

python scripts/snowflake_query_helper.py grant --role analyst_role --database analytics --schemas public,staging --privileges SELECT,USAGE

---

SQL 最佳实践

命名与风格

  • 所有标识符使用 snake_case。避免使用双引号标识符 —— 它们会强制执行区分大小写的名称,且每次使用都必须加引号。
  • 优先使用 CTE (WITH 子句) 而非嵌套子查询。
  • 使用 CREATE OR REPLACE 实现幂等 DDL。
  • 使用明确的列列表 —— 生产环境中严禁使用 SELECT *。Snowflake 的列式存储仅扫描被引用的列,因此明确列出可减少 I/O。

存储过程 —— 冒号前缀规则

在 SQL 存储过程(BEGIN...END 块)中,变量和参数在 SQL 语句内部必须使用冒号 : 前缀。否则,Snowflake 会将其视为列标识符并抛出 "invalid identifier" 错误。

sql
-- 错误:缺少冒号前缀
SELECT name INTO result FROM users WHERE id = p_id;

-- 正确:变量和参数均使用冒号前缀
SELECT name INTO :result FROM users WHERE id = :p_id;

此规则适用于在 SELECT, INSERT, UPDATE, DELETE 或 MERGE 中使用的 DECLARE 变量、LET 变量和过程参数。

半结构化数据

  • 使用 VARIANT, OBJECT, ARRAY 处理 JSON/Avro/Parquet/ORC。
  • 访问嵌套字段:src:customer.name::STRING。务必使用 ::TYPE 进行类型转换。
  • VARIANT null 与 SQL NULL:JSON 的 null 被存储为字符串 "null"。加载时请使用 STRIP_NULL_VALUE = TRUE
  • 展开数组:SELECT f.value:name::STRING FROM my_table, LATERAL FLATTEN(input => src:items) f;

使用 MERGE 实现 Upserts

sql
MERGE INTO target t USING source s ON t.id = s.id
WHEN MATCHED THEN UPDATE SET t.name = s.name, t.updated_at = CURRENT_TIMESTAMP()
WHEN NOT MATCHED THEN INSERT (id, name, updated_at) VALUES (s.id, s.name, CURRENT_TIMESTAMP());

> 深入了解 SQL 模式和反模式,请参阅 references/snowflake_sql_and_pipelines.md

---

数据管道

方案选择

| 方案 | 使用场景 |
|----------|-------------|
| Dynamic Tables | 声明式转换。默认选择。 定义查询,由 Snowflake 处理刷新。 |
| Streams + Tasks | 命令式 CDC。适用于过程化逻辑、调用存储过程、复杂的分支逻辑。 |
| Snowpipe | 从云存储(S3, GCS, Azure)持续加载文件。 |

Dynamic Tables

sql
CREATE OR REPLACE DYNAMIC TABLE cleaned_even
sql
TARGET_LAG = '5 minutes'
    WAREHOUSE = transform_wh
    AS
    SELECT event_id, event_type, user_id, event_timestamp
    FROM raw_events
    WHERE event_type IS NOT NULL;

关键规则:

  • 渐进式设置 TARGET_LAG:DAG 上游设置较紧,下游设置较松。

  • 增量 DT(动态表)不能依赖于全量刷新 DT。

  • SELECT * 在上游架构变更时会失效 —— 请使用明确的列清单。

  • 在 DAG 中,两个动态表之间不能插入视图。

Streams 和 Tasks

sql
CREATE OR REPLACE STREAM raw_stream ON TABLE raw_events;

CREATE OR REPLACE TASK process_events
WAREHOUSE = transform_wh
SCHEDULE = 'USING CRON 0 */1 * * * America/Los_Angeles'
WHEN SYSTEM$STREAM_HAS_DATA('raw_stream')
AS INSERT INTO cleaned_events SELECT ... FROM raw_stream;

-- Task 默认处于 SUSPENDED(暂停)状态。你必须手动恢复。
ALTER TASK process_events RESUME;

> 有关 DT 调试查询和 Snowpipe 模式,请参阅 references/snowflake_sql_and_pipelines.md

---

Cortex AI

函数参考

| 函数 | 用途 |
|----------|---------|
| AI_COMPLETE | LLM 补全(文本、图像、文档) |
| AI_CLASSIFY | 将文本分类(最多 500 个标签) |
| AI_FILTER | 对文本或图像进行布尔过滤 |
| AI_EXTRACT | 从文本/图像/文档中提取结构化数据 |
| AI_SENTIMENT | 情感得分 (-1 到 1) |
| AI_PARSE_DOCUMENT | 从文档中进行 OCR 或布局提取 |
| AI_REDACT | 从文本中删除 PII(个人可识别信息) |

已弃用的名称(请勿使用): COMPLETE, CLASSIFY_TEXT, EXTRACT_ANSWER, PARSE_DOCUMENT, SUMMARIZE, TRANSLATE, SENTIMENT, EMBED_TEXT_768

TO_FILE -- 常见陷阱

Stage 路径和文件名是分开的参数:

sql
-- 错误:使用单个合并参数
TO_FILE('@stage/file.pdf')

-- 正确:使用两个参数
TO_FILE('@db.schema.mystage', 'invoice.pdf')

Cortex Agents

Agent 规范使用 JSON 结构,包含顶层键:models, instructions, tools, tool_resources

  • 使用 $spec$ 分隔符(而非 $$)。
  • models 必须是一个对象,不能是数组。
  • tool_resources 是一个独立的顶层键,而不是嵌套在 tools 内部。
  • 工具描述(Tool descriptions)是决定 Agent 质量的最关键因素。

> 有关完整的 Agent 规范示例和 Cortex Search 模式,请参阅 references/cortex_ai_and_agents.md

---

Snowpark Python

python
from snowflake.snowpark import Session
import os

session = Session.builder.configs({
"account": os.environ["SNOWFLAKE_ACCOUNT"],
"user": os.environ["SNOWFLAKE_USER"],
"password": os.environ["SNOWFLAKE_PASSWORD"],
"role": "my_role", "warehouse": "my_wh",
"database": "my_db", "schema": "my_schema"
}).create()

  • 切勿硬编码凭据。请使用环境变量或密钥对认证。
  • DataFrame 是惰性的 —— 仅在调用 collect() / show() 时执行。
  • 不要对大型 DataFrame 调用 collect()。请使用 DataFrame 操作在服务端进行处理。
  • 对于批处理和 ML 工作负载,请使用 向量化 UDF(速度快 10-100 倍)。

dbt on Snowflake

sql
-- 动态表物化(用于流式/近实时数据集):
{{ config(materialized='dynamic_table', snowflake_warehouse='transforming', target_lag='1 hour') }}

-- 增量物化(用于大型事实表):
{{ config(materialized='incremental', unique_key='event_id') }}

-- Snowflake 特定配置(可与任何物化方式结合使用):
{{ config(transient=true, copy_grants=true, query_tag='team_daily') }}

  • 不要在没有 {% if is_incremental() %} 的情况下使用 {{ this }}
  • 对于流式或近实时数据集市,请使用 dynamic_table 物化。

性能优化

  • 集群键 (Cluster keys):仅用于 TB 级以上的大表。应用于 WHERE / JOIN / GROUP BY 列。
  • 搜索优化 (Search Optimization)ALTER TABLE t ADD SEARCH OPTIMIZATION ON EQUALITY(col);
  • 仓库规模 (Warehouse sizing):从 X-Small 开始,按需扩容。设置 AUTO_SUSPEND = 60AUTO_RESUME = TRUE
  • 按工作负载分离仓库(加载、转换、查询)。

安全性

  • 遵循最小权限 RBAC 原则。使用数据库角色进行对象级授权。
  • 定期审计 ACCOUNTADMIN:SHOW GRANTS OF ROLE ACCOUNTADMIN;
  • 使用网络策略进行 IP 白名单管理。
  • 对 PII 列使用掩码策略 (Masking Policies),对多租户隔离使用行访问策略 (Row Access Policies)。

---

主动触发提醒

在上下文中发现以下问题时,请主动指出:

  • SQL 存储过程缺失冒号前缀 —— 立即标记,这会导致运行时出现 "invalid identifier" 错误。
  • **动态表 (Dynamic Tables) 中使用 SELECT * —— 标记为架构变更的“定时炸弹”。
  • 过时的 Cortex 函数名称**(如 CLASSIFY_TEXT, SUMMARIZE 等) —— 建议使用当前的 AI_* 等效函数。
  • 任务创建后未恢复 (Resume) —— 提醒任务在创建时默认处于 SUSPENDED 状态。
  • Snowpark 代码中硬编码凭据 —— 标记为安全风险。

---

常见错误

| 错误 | 原因 | 解决方法 |
|-------|-------|-----|
| "Object does not exist" | 数据库/架构上下文错误或缺失权限 | 使用全限定名 (db.schema.table),检查权限 |
| 存储过程中 "Invalid identifier" | 变量缺失冒号前缀 | 在 SQL 语句中使用 :variable_name |
| "Numeric value not recognized" | VARIANT 字段未转换类型 | 显式转换:src:field::NUMBER(10,2) |
| 任务未运行 | 创建后忘记恢复 | ALTER TASK task_name RESUME; |
| DT 刷新失败 | 上游架构变更或未启用变更跟踪 | 使用显式列名,验证 change tracking |
| TO_FILE 错误 | 将组合路径作为单个参数传递 | 分成两个参数:TO_FILE('@stage', 'file.pdf') |

---

实践工作流

工作流 1:构建报表流水线 (30 分钟)

1. 暂存原始数据:创建指向 S3/GCS/Azure 的外部暂存区 (External Stage),配置 Snowpipe 实现自动摄入。
2. 使用动态表清洗:创建 TARGET_LAG = '5 minutes' 的 DT,用于过滤空值、转换类型和去重。
3. 使用下游 DT 聚合:创建第二个 DT,将清洗后的数据与维度表关联并计算指标。
4. 通过安全视图公开:为 BI 工具/API 层创建 SECURE VIEW
5. 授予权限:使用 snowflake_query_helper.py grant 生成 RBAC 语句。

工作流 2:为现有数据添加 AI 分类

1. 确定列:找到需要分类的文本列(例如:支持工单、评论)。
2. 使用 AI_CLASSIFY 测试SELECT AI_CLASSIFY(text_col, ['bug', 'feature', 'question']) FROM table LIMIT 10;
3. 创建增强 DT:创建动态表,自动对新行运行 AI_CLASSIFY
4. 监控成本:Cortex AI 按 Token 计费 —— 在全表运行前先进行采样测试。

工作流 3:调试失败的流水线

1. 检查任务历史SELECT * FROM TABLE(INFORMATION_SCHEMA.TASK_HISTORY()) WHERE STATE = 'FAILED' ORDER BY SCHEDULED_TIME DESC;
2. 检查 DT 刷新情况SELECT * FROM TABLE(INFORMATION_SCHEMA.DYNAMIC_TABLE_REFRESH_HISTORY('my_dt')) ORDER BY REFRESH_END_TIME DESC;
3. 检查 Stream 过时情况SHOW STREAMS; -- 检查 stale_after 列
4. 查阅故障排除参考指南
有关特定错误的修复方案,请参阅 references/troubleshooting.md

---

反模式 (Anti-Patterns)

| 反模式 | 失败原因 | 更好的做法 |
|---|---|---|
| 在动态表中执行 SELECT * | 上游架构变更会导致动态表静默失效 | 使用明确的列清单 |
| 存储过程中缺少冒号前缀 | 运行时出现 "Invalid identifier" 错误 | 在 SQL 块中始终使用 :variable_name |
| 所有工作负载共用单个仓库 | 加载、转换和查询之间存在资源竞争 | 按工作负载类型分离仓库 |
| 在 Snowpark 中硬编码凭据 | 安全风险,且在 CI/CD 中失效 | 使用 os.environ[] 或密钥对认证 |
| 对大型 DataFrame 执行 collect() | 将整个结果集拉取到客户端内存 | 使用 DataFrame 操作在服务端处理 |
| 使用嵌套子查询而非 CTE | 可读性差,难以调试,且 Snowflake 对 CTE 的优化更好 | 使用 WITH 子句 |
| 使用已弃用的 Cortex 函数 | CLASSIFY_TEXTSUMMARIZE 等将被移除 | 使用 AI_CLASSIFYAI_COMPLETE 等 |
| 任务未配置 WHEN SYSTEM$STREAM_HAS_DATA | 即使没有新数据,任务仍按计划运行,浪费额度 | 为流驱动的任务添加 WHEN 子句 |
| 使用双引号标识符 | 强制所有查询必须区分大小写 | 使用不加引号的 snake_case 标识符 |

---

交叉引用

| 技能 | 关系 |
|-------|-------------|
| engineering/sql-database-assistant | 通用 SQL 模式 —— 用于非 Snowflake 数据库 |
| engineering/database-designer | 架构设计 —— 用于 Snowflake 实现前的数据建模 |
| engineering-team/senior-data-engineer | 更广泛的数据工程 —— 流水线、Spark、Airflow、数据质量 |
| engineering-team/senior-data-scientist | 分析与机器学习 —— 与 Snowpark 结合用于特征工程 |
| engineering-team/senior-devops | Snowflake 部署的 CI/CD (Terraform, GitHub Actions) |

---

参考文档

| 文档 | 内容 |
|----------|----------|
| references/snowflake_sql_and_pipelines.md | SQL 模式、MERGE 模板、动态表调试、Snowpipe、反模式 |
| references/cortex_ai_and_agents.md | Cortex AI 函数、Agent 规范结构、Cortex Search、Snowpark |
| references/troubleshooting.md | 错误参考、调试查询、常见修复方案 |