Snowflake 开发
Snowflake 开发
涵盖 Snowflake SQL、数据管道、Cortex AI 和 Snowpark Python 开发。包括冒号前缀规则、半结构化数据、MERGE upserts、Dynamic Tables、Streams+Tasks、Cortex AI 函数、Agent 规范、性能调优和安全加固。
> 最初由 James Cha-Earley 贡献 —— 由 claude-skills 团队增强并集成。
快速上手
# 生成 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" 错误。
-- 错误:缺少冒号前缀
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
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
CREATE OR REPLACE DYNAMIC TABLE cleaned_evenTARGET_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
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 路径和文件名是分开的参数:
-- 错误:使用单个合并参数
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
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
-- 动态表物化(用于流式/近实时数据集):
{{ 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 = 60和AUTO_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_TEXT、SUMMARIZE 等将被移除 | 使用 AI_CLASSIFY、AI_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 | 错误参考、调试查询、常见修复方案 |