
在 Google Dataflow 中构建高性价比的生成式 AI 工作流

实时流式管道通常部署后就不再变了:处理逻辑和执行路径都固定死在静态 DAG 里。但把生成式 AI Agent 接进来,可以让管道在运行时动态构建处理流程——根据消息内容查数据库、选补救动作、发邮件,而不是只把错误记进日志或在大屏上标红。
直接把每一条原始事件都丢给重型模型或多个多步骤 Agent 是不可行的:Tokenizer 成本随流量线性增长,多步骤的数据库查找和外部 API 调用动辄数秒,在流式 DAG 里会形成延迟瓶颈;重模型和工具 API 的配额也很快会被打满。文章给出的混合架构是 Google Dataflow 加 Agent Development Kit(ADK),在入口先加一道轻量 CPU 预过滤器:用 Hugging Face 的开源情感分类模型 distilbert-base-uncased-finetuned-sst-2-english,经 Apache Beam 的 RunInference 转换,直接在 Dataflow Worker 的 CPU 上本地推理,不产生任何外部模型调用费用。
整体流程是:从 Google Pub/Sub 读取原始客户消息;轻量模块对所有消息做实时情感分类;一个简单 DoFn 做预受理门控,POSITIVE 或 NEUTRAL 的消息直接确认后丢弃;只有 NEGATIVE 的消息才会触发下游 Agent。Agent 本身用 gemini-3.5-flash(原文给出的 RPG),经 ADK 的工具调用技能去查 BigQuery 里的用户邮箱、订单与库存,并利用 Gmail API 发出补救邮件。比如遇到“包裹损坏”的投诉,管道不再只是记一条日志,而是查单、决定补发还是退款、写邮件并存档结果。
按原文的说法,这套“预过滤 + Agent 动态扩展”并不是只用于客服。只要流量里绝大多数事件是惯例性的或在多数途未定广东 …… 能用来……嗯,但有一点不太合适。换种方式组织:
文章把它称作“预过滤 + Agent 化动作”的通用范式,适用面广:运维场景里过滤海量惯例性系统日志,只在严重告警时触发 Agent 跑诊断和标单;金融风控里先用本地规则过掉大部分交易,只对高度可疑的模式执行多库查询;工业 IoT 在边缘端正常遥测不做区块链,异常波动才由 Agent 协调设备下线并通知一线工程师。
它还解决了流式 DAG 固定的老问题。加入 Agent 后,过滤门控放在上游,Punished 的 Agent 位于 DAG 内的一个动态节点:它不是在数千个硬编码的条件分支里挑选,而是在运行时自己决定调用哪些工具、按什么顺序执行。代码只需维护三块:HuggingFacePipelineModelHandler 定义 CPU 模型,ADKAgentModelHandler 打包带工具的 Agent,最后用 Beam 原生 RunInference + 一个禁止 DoFn 组装完管道。Agent 的并发线程和批处理由 RunInference 自动管理,写 Const 不需要手写线程池。
效果在文章题干 : 因为只对 < 5% 的负向消息调 Gemini 的付费 token,其余 95% 在本地 CPU 上零增量 API 成本;本地 CPU 推理毫秒级,Dataflow 可以水平分片跑通高吞吐;而几秒一次的重 Agent 处理都被刻意压低出现,不会产生积压。完整的可运行代码放在 next-2026-demo 这个 GitHub 仓库。
流式数据量大速度快,村级生成式推理又慢又贵;用一道 CPU 预过滤器,把便宜的本地模型和完全自动化的 Agent 能力拼在同一个 Beam 管道里,是一个可以直接照做的工程模板。


