Snowflake AI函数消除80%自定义管道代码
译文 AI 逐段翻译
用内联 SQL 取代定制的 Python 脚本、外部 API 和编排胶水代码

免责声明:我是Snowflake 首席技术架构师拥有30多年的数据战略、架构和开发经验。此处表达的观点仅代表我个人,不一定反映现任、前任或未来雇主的观点
引言
数据工程师花费大量时间构建的管道,实际上并非关乎数据移动,而是文本增强。对工单进行分类,从电子邮件中提取实体,对评论进行情感评分,翻译反馈,在共享前对 PII 进行脱敏。
这些“简单”任务中的每一个传统上都需要:
- 带有机器学习库依赖的 Python/Spark 作业
- 外部 API 调用(OpenAI、AWS Comprehend、Google NLP)
- 数据从仓库导出到计算层
- 错误处理、重试、速率限制和分页
- 将结果摄取回仓库
- 编排(Airflow、Prefect、Dagster)用于调度和监控
- API 密钥的机密管理
- 需要维护和扩展的独立基础设施
Snowflake AI 函数将整个技术栈折叠为一条 SQL 语句。数据永不离开。没有外部 API。无需管理机密。没有编排层。它只是工作。
管道税:工程师实际构建的东西
一个典型的数据工程团队维护着许多类似这样的“增强”管道:
之前:传统管道
- Airflow DAG 每小时触发
- Python 任务查询 Snowflake 获取新行(SELECT … WHERE processed_at IS NULL)
- 对于每个批次:调用 OpenAI/Comprehend API 处理文本
- 解析 JSON 响应,处理错误,重试失败
- 将结果写回 Snowflake(MERGE 或 INSERT)
- 标记行已处理
- 失败时告警

该管道需要:Airflow、Python、OpenAI 密钥、网络出口、错误处理、幂等逻辑和监控。
之后:一条 SQL 语句
CREATE OR REPLACE DYNAMIC TABLE ENRICHED_TICKETS
TARGET_LAG = '1 hour'
WAREHOUSE = ADMIN_DB_WH
AS
SELECT
ticket_id,
ticket_text,
customer_id,
AI_CLASSIFY(
ticket_text,
['billing', 'technical', 'account', 'general', 'product_feedback']
):labels[0]::VARCHAR AS category,
AI_SENTIMENT(ticket_text):categories[0]:sentiment::VARCHAR AS sentiment,
TO_JSON(AI_EXTRACT(
ticket_text,
['customer_action', 'urgency_level', 'competitor_mentioned']
)) AS extracted_entities,
created_at
FROM support_tickets;
这取代了整个管道。没有 Airflow。没有 Python。没有 API 密钥。没有出口。没有重试逻辑。没有监控基础设施。
注意:这段确切的代码已部署并在 AGUHA_AI.DEMO 上运行,8 行数据已完全增强。动态表处于活动状态,采用增量刷新模式。

你可以删除的七种管道
1. 文本分类管道
用例:对支持工单进行分类、路由电子邮件、标记文档
之前:Python + scikit-learn 模型或 OpenAI API + Airflow DAG + 模型注册表
之后(已验证):
AI_CLASSIFY(ticket_text, ['billing','technical','account','general','product_feedback']):labels[0]::VARCHAR影响:消除了模型训练、部署、版本管理和推理基础设施。标签在 SQL 中定义,而不是模型工件。
2. 实体提取管道
用例:从非结构化文本中提取姓名、日期、金额、产品提及
之前:spaCy/Hugging Face NER 模型 + Python 服务 + 结果解析 + Snowflake MERGE
之后(已验证):
TO_JSON(AI_EXTRACT(ticket_text, ['customer_action','urgency_level','competitor_mentioned']))
- Access fields: AI_EXTRACT(…):response:urgency_level::VARCHAR影响:无需训练或部署模型。字段在查询时指定。无需重新训练即可添加新字段。
3. PII 检测与脱敏管道
用例:在共享或分析之前查找并编辑个人数据
之前:Presidio/AWS Macie + 扫描作业 + 脱敏视图 + 调度
之后(已验证):
AI_REDACT(feedback_text)
- Masks: [email protected] -> [EMAIL], 555–0123 -> [PHONE_NUMBER], Sarah Johnson -> [NAME]影响:一个函数调用取代了整个 PII 扫描基础设施。在视图中使用可实现动态脱敏。
4. 情感分析管道
用例:对客户反馈、NPS 评论、社交媒体提及进行评分
之前:VADER/TextBlob 或 API 调用 + 批处理 + 结果聚合
之后(已验证):
AI_SENTIMENT(review_text):categories[0]:sentiment::VARCHAR
- Returns: 'positive', 'negative', 'mixed', 'neutral'
- NOTE: Does NOT return a float score. Returns categorical labels.影响:内联在任何查询中。无需批处理作业。与 GROUP BY 结合可立即获得按产品的情感仪表板。
5. 翻译管道
用例:翻译客户内容、支持工单、产品描述
之前:Google Translate API + 速率限制 + 批处理 + 成本跟踪
之后(已验证):
AI_TRANSLATE(feedback_text, 'fr', 'en') - French to English
AI_TRANSLATE(feedback_text, 'de', 'en') - German to English影响:无需 API 密钥、无速率限制、无出口。在查询时翻译或使用动态表物化。
6. 响应生成管道
用例:草拟客户支持响应、总结反馈、生成报告
之前:OpenAI API + 提示管理 + 响应缓存 + 成本控制
之后(已验证):
AI_COMPLETE(
'llama3.1-8b',
'Write a brief, professional 2-sentence customer support response: ' || ticket_text
)注意:模型可用性因账户而异。llama3.1–8b 确认可用。claude-3–5-sonnet 和 snowflake-arctic 并非所有账户都可用。
影响:无需 API 密钥、无需提示版本管理基础设施、无需管理外部 LLM 服务。
7. 聚合洞察管道
用例:总结数百条反馈条目中的主题、发现模式
之前:自定义 Python 脚本 + LLM API + 分块逻辑 + 结果聚合
之后(已验证):
SELECT AI_AGG(feedback_text, 'What are the top 3 themes across this feedback?') AS themes
FROM customer_feedback;
- Returns structured multi-paragraph analysis across all rows影响:一个 SQL 函数取代了复杂的分块 + 摘要管道。无需管理令牌限制。
架构转变
之前:提取-增强-加载(EEL)
Snowflake -> 提取到 Python -> 调用外部 AI -> 解析结果 -> 加载回 Snowflake
问题:数据出口、API 故障、过时结果、复杂编排、多故障点。
之后:原地增强(EiP)
原始表 -> 带有 AI 函数的动态表 -> 增强表
好处:零出口、原子事务、自动重试、增量处理、单一系统。
这对团队意味着什么
- 数据工程师不再是 API 管道工,而是专注于数据建模
- 简单的分类/提取任务不再需要 ML 运维团队
- 分析工程师可以将 AI 增强添加到他们的 dbt 模型中
- 上线时间从数周缩短到数分钟
- 成本从基础设施 + API 费用转移到 Snowflake 积分(单一账单)
- 安全团队满意:无数据出口、无需轮换外部 API 密钥
使用 AI 函数的生产架构

第 1 层:原始摄取(不变)
Snowpipe、Snowpipe Streaming 或 COPY INTO 按原样加载原始数据。不在摄取时进行增强——保持简单和快速。
第 2 层:AI 增强(动态表)
带有 AI 函数的动态表将原始文本转换为结构化、分类、情感评分、实体提取的数据。刷新是自动且增量的。
已验证的生产动态表:
CREATE OR REPLACE DYNAMIC TABLE ENRICHED_TICKETS
TARGET_LAG = '1 hour'
WAREHOUSE = ADMIN_DB_WH
AS
SELECT
ticket_id, ticket_text, customer_id,
AI_CLASSIFY(ticket_text, ['billing','technical','account','general','product_feedback']):labels[0]::VARCHAR AS category,
AI_SENTIMENT(ticket_text):categories[0]:sentiment::VARCHAR AS sentiment,
TO_JSON(AI_EXTRACT(ticket_text, ['customer_action','urgency_level','competitor_mentioned'])) AS extracted_entities,
created_at
FROM support_tickets;第 3 层:消费(视图、仪表板、告警)
下游消费者从增强后的动态表中读取数据。分析师查询结构化字段。仪表盘可视化情感趋势。警报根据紧急程度触发。
增量处理:只为新数据付费
带有AI函数的动态表默认是增量的。当源表中有新行到达时,只有这些行会在下次刷新时通过AI函数。历史数据不会被重新处理。
这是关键的成本控制机制。一个有1000万历史行和每小时1000新行的表,每次刷新周期只处理1000行——而不是1000万行。
AI函数不能替代什么
坦诚地说明局限性:
- 需要领域特定准确性的自定义微调模型
- 延迟低于100毫秒的实时推理(AI函数是批处理导向的)
- 一次处理视频+音频+文本的多模态流水线
- 需要写入外部系统(CRM、工单系统)的流水线
- 需要模型可解释性或置信度分数的任务
- AI_FILTER(某些账户存在语法错误,导致无法使用)
- 除llama3.1–8b之外的模型(可用性因账户/地区而异)
迁移手册:替换现有流水线
步骤1:盘点你的增强流水线
列出所有调用外部AI/NLP API或运行ML模型进行分类、提取、情感分析、翻译或PII掩码的流水线。
步骤2:为每个流水线评估AI函数的适用性
高适用性:文本分类、实体提取、情感分析、翻译、PII掩码 中适用性:文档解析、摘要、响应生成 低适用性:自定义模型、实时推理、多模态
步骤3:从最简单、最高吞吐量的流水线开始
选择处理行数最多且逻辑最简单的流水线。这样能带来最大的投资回报率和最小的风险。
步骤4:将替换方案构建为动态表
编写包含AI函数的SELECT语句。设置TARGET_LAG以匹配你当前的SLA。先在样本上测试(LIMIT 10),然后在完整数据上启用。
步骤5:并行运行一个周期
比较旧流水线和新动态表的结果。验证准确性是否等同。检查信用成本与旧基础设施成本的对比。
步骤6:停用旧流水线
关闭Airflow DAG。删除Python代码。取消API密钥。删除Lambda函数。更新文档。
真实示例:
配置(已部署并运行)
- 3个源表 support_tickets — 8行,客户支持文本 product_reviews — 8行,带评分的产品反馈 customer_feedback — 6行,带PII的多语言反馈
- 1个动态表(自动增强) ENRICHED_TICKETS — 结合了AI_CLASSIFY + AI_SENTIMENT + AI_EXTRACT — 状态:ACTIVE,模式:INCREMENTAL,8/8行已填充
- 1个视图(PII保护) SAFE_FEEDBACK — AI_REDACT在读取时动态掩码所有PII
这替代了什么
如果没有AI函数,要实现相同的结果需要:
- 3个独立的ML模型(分类、情感分析、NER)或3个API集成
- 1个PII扫描服务(Presidio、AWS Macie)
- 1个编排工具(Airflow)和4个DAG
- 用于API调用、错误处理和结果解析的Python代码
- API密钥的密钥管理
- 每个流水线的监控和告警
- 大致基础设施:跨3–4个服务的6–8个组件
使用AI函数:
- 1个动态表 + 1个视图 = 完整的增强流水线
- 0个外部服务
- 0个API密钥
- 0个编排
- 0个持续维护

结论
AI函数并非取代数据工程——它们消除了那些消耗了增强流水线开发时间60–80%的无法区分的繁重工作。分类逻辑、情感分析、实体提取、翻译、PII掩码——所有这些都变成了SQL关注点,而非基础设施关注点。
这种转变既是技术性的,也是理念性的:增强变成了一个SQL列,而不是一个微服务。数据工程师可以专注于数据建模、质量和治理,而不是管理API管道。
从一个流水线开始。用动态表替换它。衡量节省的成本。然后做下一个。你每添加一个AI函数,就多一行SQL,少一个需要维护的服务。
Snowflake AI函数如何消除80%的自定义流水线代码 最初发表在 Snowflake Builders Blog: Data Engineers, App Developers, AI, & Data Science 在 Medium 上,人们在那里通过强调和回应这个故事来继续对话。