返回
RSS Snowflake Engineering (Medium) 精选 发布 2026-08-13 22:01 收录于 08-14 60

Snowflake AI函数消除80%自定义管道代码

Snowflake AI Functions将文本分类、实体提取、情感分析等数据丰富任务从传统的Python脚本、外部API和编排工作流中抽离,合并为一条SQL语句。据称可减少80%的自定义管道代码,数据无需离开仓库,也无需管理API密钥和编排层。数据工程师可直接在Snowflake内完成文本增强,大幅简化管道架构。
推荐理由:数据从业者可了解平台内置AI函数如何简化数据管道,减少外部依赖与运维成本。
平台AI化Snowflake

译文 AI 逐段翻译

用内联 SQL 取代定制的 Python 脚本、外部 API 和编排胶水代码

免责声明:我是Snowflake 首席技术架构师拥有30多年的数据战略、架构和开发经验。此处表达的观点仅代表我个人,不一定反映现任、前任或未来雇主的观点

引言

数据工程师花费大量时间构建的管道,实际上并非关乎数据移动,而是文本增强。对工单进行分类,从电子邮件中提取实体,对评论进行情感评分,翻译反馈,在共享前对 PII 进行脱敏。

这些“简单”任务中的每一个传统上都需要:

  • 带有机器学习库依赖的 Python/Spark 作业
  • 外部 API 调用(OpenAI、AWS Comprehend、Google NLP)
  • 数据从仓库导出到计算层
  • 错误处理、重试、速率限制和分页
  • 将结果摄取回仓库
  • 编排(Airflow、Prefect、Dagster)用于调度和监控
  • API 密钥的机密管理
  • 需要维护和扩展的独立基础设施

Snowflake AI 函数将整个技术栈折叠为一条 SQL 语句。数据永不离开。没有外部 API。无需管理机密。没有编排层。它只是工作。

管道税:工程师实际构建的东西

一个典型的数据工程团队维护着许多类似这样的“增强”管道:

之前:传统管道

  1. Airflow DAG 每小时触发
  2. Python 任务查询 Snowflake 获取新行(SELECT … WHERE processed_at IS NULL)
  3. 对于每个批次:调用 OpenAI/Comprehend API 处理文本
  4. 解析 JSON 响应,处理错误,重试失败
  5. 将结果写回 Snowflake(MERGE 或 INSERT)
  6. 标记行已处理
  7. 失败时告警

该管道需要: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 上,人们在那里通过强调和回应这个故事来继续对话。

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