Google Dataflow 结合 ADK 构建经济高效的生成式 AI 流式管道
DataHot 速览
Google Cloud 博客介绍了在 Dataflow 中集成生成式 AI 代理的混合流式管道模式。通过在上游使用轻量级 CPU 机器学习模型过滤事件,仅将复杂案例路由至下游 ADK 代理,实现了流式处理的动态分支,同时显著降低成本与延迟。文中以客户支持场景为例,说明代理可动态查询订单数据库、制定补救措施并记录结果。该模式适用于高吞吐流中大多数事件为常规场景的通用场景。
为什么值得关注:数据从业者可借鉴其在流式数据处理中结合轻量预过滤与代理动态决策的架构,兼顾成本、延迟与复杂事件处理能力。
本文目录 13 节
译文
AI 逐段翻译实时流式管道是现代企业的运营支柱,持续处理从客户支持互动到交易日志的一切内容。传统上,流式DAG是静态的;一旦部署,其处理逻辑和执行路径便固定不变。然而,通过集成生成式AI代理,我们可以超越静态逻辑,实现自适应执行。这使得流式工作流能够在运行时根据数据内容动态构建计划、查询数据库并触发自定义修复路径。
例如,当客户发送关于损坏订单的愤怒消息时,管道不应仅仅记录错误或在仪表板上标记。它应该查找包含客户订单和库存记录的数据库,决定修复操作(如发送替换品或退款),给客户发送邮件,并记录最终解决情况。
然而,流式系统在执行生成式AI工作流时面临一个基本的工程障碍:规模、延迟和成本。将每个原始事件直接发送到重型模型或配备外部数据库和邮件工具的多步代理,成本过高,引入高延迟,并迅速耗尽API速率限制。
此模式通过结合Google Dataflow(Google Cloud的完全托管、无服务器执行服务,用于Apache Beam)和代理开发套件(ADK)来构建混合流式管道,从而解决规模和复杂性挑战。通过使用轻量级、CPU绑定的机器学习模型在上游过滤和限定事件,我们保持管道的高度成本效益,仅将复杂案例路由到下游代理。在那里,代理动态决定采取什么行动,将动态分支引入流中,而无需在管道的静态DAG中硬编码数千个条件步骤。
高容量流的通用蓝图
虽然我们下面使用客户支持分类场景,但这种预过滤+代理行动模式是一种通用范式。它适用于任何具有高容量(>9X%)常规事件且只有少数需要复杂上下文推理的流。
- IT运营与DevOps:在CPU上过滤数百万条常规系统日志,仅在标记严重异常时触发代理运行诊断并创建Bug工单。
- 金融欺诈分类:通过轻量级本地规则传递数百万笔交易,仅对高度可疑的模式调用代理执行多数据库查找工具。
- 工业物联网:在边缘监控正常遥测,并将异常尖峰路由到代理以协调设备停机并向现场工程师发送邮件。
架构:为什么预过滤流式事件?
在高吞吐量流中,绝大多数消息不需要复杂推理或修复。它们可能是积极反馈、中性询问或简单查询。
将每个事件路由到重型LLM工作流会产生三个主要瓶颈:
- API成本:前沿模型按令牌收费。在高吞吐量下,成本随流容量线性增长。
- 延迟:多步工作流(涉及数据库查找和外部API调用)需要数秒,在流式DAG中造成瓶颈。
- 配额:外部API有严格的速率限制,流式工作者很容易耗尽。
为防止这种情况,我们在Apache Beam/Dataflow中构建预过滤管道:
管道流程
- 摄取:从Google Pub/Sub读取原始客户消息。
- 轻量级情感分类器(CPU):使用Apache Beam的
distilbert-base-uncased-finetuned-sst-2-english通过RunInference转换,在Dataflow工作者的CPU上本地运行所有消息,避免外部API成本。 - 预资格门:一个简单的
DoFn过滤流。带有POSITIVE或NEUTRAL情感的消息被确认并丢弃。 - 自动化修复(ADK):当且仅当消息被分类为
NEGATIVE时,我们触发由gemini-3.5-flash支持的生成式AI代理,使用ADKAgentModelHandler。代理使用工具在BigQuery中查找用户、获取订单、选择修复计划,并通过Gmail API发送通知邮件。
自适应执行:使Beam DAG动态化
在传统流式架构中,管道的有向无环图(DAG)是刚性的。一旦部署到Dataflow,转换序列就固定了。如果需要处理新型警报或更改特定事件的路由方式,必须修改、测试并重新部署整个管道。
通过在情感预过滤器之后放置生成式AI代理,我们在静态DAG中引入了一个动态、自适应的节点。
对于95%的积极或中性记录,管道沿着快速、静态路径运行。但当过滤器门控到消极记录时,代理评估负载并在运行时动态选择正确的API工具序列(如数据库查询、库存检查或邮件通知)。这允许管道动态执行复杂决策树,无需在静态Apache Beam代码中构建和维护数千个硬编码条件分支。
实现管道
以下是使用Google代理开发套件(ADK)和RunInference框架在Apache Beam中的示例实现。
1. 定义轻量级情感模型
我们使用HuggingFacePipelineModelHandler定义上游CPU模型。该模型将情感分类为POSITIVE、NEUTRAL或NEGATIVE在工作者实例上。
加载中...
2. 构建重型ADK代理
ADK代理充当我们的修复助手。我们为其配备三个工具:
-
lookup_user:查询BigQuery获取客户电子邮件。 -
lookup_orders:查询BigQuery获取客户订单和当前产品库存。 -
send_email:使用Gmail API向客户发送修复邮件。
加载中...
我们配置LlmAgent并将其封装在ADKAgentModelHandler中:
加载中...
3. 组装Dataflow DAG
整个流水线被清晰地声明。上游情感推理直接馈入过滤步骤(FilterNegativeADK),然后有条件地执行下游ADKInference:
加载中...
成本与性能优势
通过引入此过滤步骤,我们获得了重大的工程和运营优势:
1. 显著降低成本
无需为100%的传入事件支付Gemini输入/输出令牌费用,我们仅为代表负面客户情绪的部分(通常小于5%的消息)付费。其余95%在CPU实例上本地分类,无需额外API成本。
2. 高流式吞吐量
Dataflow将CPU分类工作负载分布到多个实例。由于CPU推理仅需毫秒级,流水线可水平扩展以处理高吞吐的事件流。重量级LLM代理因工具执行可能每秒处理多个请求,被谨慎调用,避免积压。
3. 原生Apache Beam集成
将代理添加到DAG中无需复杂的编排逻辑或手动线程池。使用ADKAgentModelHandler与Beam的原生RunInference转换自动处理并行工作线程、批处理和集成,保持代码库可维护且整洁。
关键要点
流式数据快速且量大,而重量级生成式AI推理缓慢且昂贵。
通过使用Google Dataflow和ADK构建预过滤流水线,您可兼得两全其美:本地CPU模型的成本与速度,以及Gemini支持的代理的深度自动化能力。
要查看完整代码库并自行部署,请参阅next-2026-demo GitHub仓库。
Apache Beam是Apache Software Foundation
的商标。发布于
这篇内容对你有用吗?
反馈只用于改善内容筛选,不等同于收藏