实时流处理擅长稳定地执行固定 DAG,生成式 AI Agent 则擅长根据上下文临时规划、查询数据并调用工具。把两者直接拼在一起并不难,难的是控制成本、延迟和外部 API 配额:如果每条消息都进入大模型和多步 Agent,高吞吐数据流很快就会变成昂贵而缓慢的任务队列。
更实用的架构是设置一道“轻模型闸门”:所有事件先由 Dataflow 工作节点上的 CPU 模型低成本分类,只有少量复杂事件才交给 Agent。这样既保留 Apache Beam 静态 DAG 的可预测性,又能在关键节点引入动态执行。
静态 DAG 中嵌入动态决策节点
以客户支持消息为例,数据路径可以拆成五步:
- 从 Pub/Sub 持续读取消息。
- 使用 Apache Beam
RunInference在工作节点本地执行轻量情感分类。 - 丢弃无需处置的消息,只保留负面事件。
- 将负面事件交给 ADK Agent。
- Agent 按上下文查询 BigQuery、检查订单与库存,并通过 Gmail API 通知客户。
从 Dataflow 的角度看,整体拓扑仍然是固定的:
Pub/Sub -> CPU 分类 -> 过滤 -> ADK Agent -> 审计输出
动态性发生在 Agent 节点内部。它可以在运行时决定调用哪些工具以及调用顺序,例如:
查询客户 -> 查询订单 -> 检查库存 -> 安排补发 -> 发送邮件
另一个事件可能采用不同路径:
查询客户 -> 查询订单 -> 判断无法补货 -> 建议退款 -> 发送邮件
因此,不必在 Beam 代码中维护成百上千个 if/elif 分支;但这不代表 Agent 可以不受约束。退款、发货、停机等有副作用的动作仍应受到权限、金额阈值、幂等和人工审批规则保护。
为什么预过滤是关键,而不只是一次模型优化
假设每秒进入 1,000 条事件,其中只有 5% 需要复杂处理。未经筛选时,Agent 每秒需要承接 1,000 次调用;加上闸门后,调用量降至约 50 次。收益同时出现在三个维度:
- 费用:大模型 Token 和工具调用成本只作用于少量候选事件。CPU 推理仍会产生 Dataflow 计算费用,但不会为每条消息支付外部大模型推理费。
- 延迟:本地分类通常远快于包含数据库和邮件工具的多步 Agent,有助于避免慢请求堵塞主数据流。
- 配额:BigQuery、Gmail 和模型 API 都可能有并发或速率限制,预过滤可以显著降低突发压力。
这个模式不局限于客服:运维日志可以先检测异常,再让 Agent 诊断并创建工单;交易流可以先由规则或小模型筛选,再执行多库反欺诈调查;工业遥测可以先过滤正常读数,仅将异常尖峰交给 Agent 协调停机和通知。
闸门的目标不是做到最聪明,而是用较低成本获得足够高的召回率。漏掉一个真正需要处置的事件,往往比多送几个误报给 Agent 更危险,因此阈值通常应偏向召回率,并通过监控持续校准。
可以这样组装 Beam 管道
下面是一个可改造的核心示例,展示 Pub/Sub、Hugging Face CPU 推理、负面消息过滤和 ADK 推理之间的连接。运行前需要安装并配置 Apache Beam、Google Cloud 客户端、Hugging Face 推理依赖和 ADK 集成包,同时将项目、区域、暂存目录和 Pub/Sub Topic 替换为自己的值。
import argparse
import json
import apache_beam as beam
from apache_beam.ml.inference.base import RunInference
from apache_beam.ml.inference.huggingface_inference import (
HuggingFacePipelineModelHandler,
)
from apache_beam.options.pipeline_options import PipelineOptions
# 按当前 ADK/Beam 集成包调整实际导入路径。
from google.adk.agents import LlmAgent
from apache_beam.ml.inference.adk import ADKAgentModelHandler
class KeepNegativeAndBuildPrompt(beam.DoFn):
def process(self, result):
prediction = result.inference
if isinstance(prediction, list):
prediction = prediction[0]
label = str(prediction.get("label", "")).upper()
if label != "NEGATIVE":
return
message = result.example
try:
event = json.loads(message)
except json.JSONDecodeError:
event = {"message": message}
yield (
"Investigate this customer event. Verify the user and order before "
"taking action. Do not invent missing data. Event: "
+ json.dumps(event, ensure_ascii=False)
)
def build_pipeline(argv=None):
parser = argparse.ArgumentParser()
parser.add_argument("--input_topic", required=True)
known, beam_args = parser.parse_known_args(argv)
options = PipelineOptions(beam_args, streaming=True, save_main_session=True)
sentiment_handler = HuggingFacePipelineModelHandler(
task="sentiment-analysis",
model="distilbert-base-uncased-finetuned-sst-2-english",
)
# lookup_user、lookup_orders、send_email 需要在生产代码中实现,
# 并加入参数校验、超时、重试、幂等和最小权限控制。
tools = [lookup_user, lookup_orders, send_email]
agent = LlmAgent(
name="remediation_agent",
model="gemini-3.5-flash",
instruction=(
"You handle verified customer remediation. Use tools to inspect "
"records before acting. Never fabricate order or inventory data. "
"Do not send duplicate messages."
),
tools=tools,
)
agent_handler = ADKAgentModelHandler(agent=agent)
with beam.Pipeline(options=options) as pipeline:
_ = (
pipeline
| "Read" >> beam.io.ReadFromPubSub(topic=known.input_topic)
| "Decode" >> beam.Map(bytes.decode, "utf-8")
| "ClassifyOnCPU" >> RunInference(sentiment_handler)
| "KeepNegative" >> beam.ParDo(KeepNegativeAndBuildPrompt())
| "RunAgent" >> RunInference(agent_handler)
| "LogResult" >> beam.Map(print)
)
if __name__ == "__main__":
build_pipeline()
需要特别注意:示例中的 distilbert-base-uncased-finetuned-sst-2-english 通常是二分类模型,输出 POSITIVE 或 NEGATIVE,并不会原生给出 NEUTRAL。如果业务确实需要正面、中性、负面三分类,应更换三分类模型或使用分数阈值定义中性区间,并用自己的数据验证效果。
开发阶段可使用 Direct Runner 验证管道结构:
python pipeline.py \
--runner=DirectRunner \
--input_topic=projects/PROJECT_ID/topics/customer-events
提交到 Dataflow 时,可以这样实践:
python pipeline.py \
--runner=DataflowRunner \
--project=PROJECT_ID \
--region=us-central1 \
--temp_location=gs://BUCKET/tmp \
--staging_location=gs://BUCKET/staging \
--input_topic=projects/PROJECT_ID/topics/customer-events \
--streaming
具体依赖版本和 ADK Handler 的导入路径应按当前 SDK 文档锁定,并先在测试项目中验证。
工具调用比 Prompt 更需要工程约束
Agent 能查询 BigQuery 和发送邮件,并不意味着应该把完整权限直接交给模型。可以这样实践:
- BigQuery 查询使用参数绑定,禁止让模型生成任意 SQL。
- 工具只返回完成决策所需的最少字段,避免泄露无关个人信息。
- 邮件工具接受结构化参数,并校验收件人与订单的归属关系。
- 用
event_id或处置编号实现幂等,防止 Dataflow 重试导致重复邮件、重复退款或重复补发。 - 给数据库和外部 API 设置超时、有限重试及死信队列。
- 将 Agent 的计划、工具参数、工具结果和最终动作写入审计日志。
- 对退款、账户冻结、设备停机等高风险动作加入规则引擎或人工批准。
例如,邮件发送工具可以先执行幂等检查,而不是收到调用就立即发送:
def send_email_once(event_id: str, to_address: str, subject: str, body: str) -> dict:
"""示意实现:生产环境应使用持久化幂等表和原子写入。"""
if remediation_already_completed(event_id):
return {"status": "skipped", "reason": "duplicate_event"}
validate_recipient_for_event(event_id, to_address)
gmail_send(to_address=to_address, subject=subject, body=body)
mark_remediation_completed(event_id)
return {"status": "sent", "event_id": event_id}
内存集合不能满足分布式幂等要求。实际部署时可将状态放入具有唯一键约束或条件写能力的持久化存储中,并考虑“邮件已发出但状态写入失败”这类部分成功场景。
上线前要观察的不是平均值
混合管道至少需要分别监控快路径和 Agent 路径:
- 每分钟输入事件数与负面事件占比;
- 轻量分类器的精确率、召回率和分数分布漂移;
- Agent 调用吞吐、P50/P95/P99 延迟与失败率;
- BigQuery、Gmail 和模型 API 的限流次数;
- Dataflow backlog、工作节点扩缩容和单事件成本;
- 工具副作用的重复率、人工撤销率和升级人工处理的比例。
如果负面事件比例从 5% 突然升到 40%,下游容量假设会立即失效。应配置并发上限、速率限制、缓冲队列和降级策略,例如暂时只创建人工工单,而不继续执行自动补发。
采用这套模式时的检查清单
这类架构适合“绝大多数事件常规、少数事件需要上下文推理”的数据流。落地前可以确认:
- 轻模型能否以可接受的漏报率筛掉 90% 以上事件?
- Agent 路径是否与主路径隔离,慢调用会不会制造无界积压?
- 每个有副作用的工具是否具备幂等、权限控制和审计记录?
- 模型或外部 API 不可用时,事件能否进入死信队列或人工队列?
- 是否用真实流量分布测算过 CPU、Token、数据库查询和邮件调用成本?
核心原则很直接:让廉价、并行、可预测的计算处理大多数事件,把昂贵而灵活的推理留给真正需要它的少数情况。Dataflow 负责稳定搬运和扩展数据,轻量模型负责守门,Agent 只在复杂事件上发挥动态规划能力。