别把每条 Kafka 消息都扔给 Agent:高吞吐 AI 流水线真正省钱的地方,是先把 90% 的普通事件挡在模型外
昨天 Google Cloud 写了一篇很“工程”的文章:用 Dataflow + Agent Development Kit 做高吞吐生成式 AI 流水线。
我看完最大的感受不是某个新 API,而是它把一个经常被忽略的架构原则说得很清楚:
高吞吐系统里,大多数事件根本不需要 Agent。
很多 AI 项目一开始都很兴奋:
Kafka 来一条消息
→ 调一次 LLM
→ Agent 判断
→ 再调用工具
PoC 每分钟几十条,看起来没问题。
一上生产:
每秒 5,000 条
成本、延迟、限流同时爆掉。
Google 给出的模式其实很朴素:
High-volume Stream
→ Cheap Pre-filter
→ 只把复杂事件交给 Agent
这篇不照抄 Dataflow 示例,我直接把这个思路改写成一个更容易放进企业 Java 栈的版本:
Kafka
→ Spring Boot Pre-filter
→ Qualified Topic
→ Agent Worker
→ Tool / Action
先算一笔账
假设客服系统每天有:
10,000,000 条事件
其中:
85% 普通通知 / 正常状态
10% 规则可以直接处理
5% 需要真正语义推理
如果 100% 都进 Agent:
10,000,000 次模型工作
如果前面加 Gate:
500,000 次
直接少一个数量级。
而且 Agent 通常不是一次模型调用。
一次复杂任务可能包含:
分类
→ Tool
→ 再推理
→ 再 Tool
→ 总结
所以实际节省不是简单的 95%。
最常见的错误:把“智能”理解成“所有数据都送模型”
例如日志告警系统:
CPU 42%
CPU 43%
CPU 41%
CPU 44%
这些正常数据为什么要进入 LLM?
IoT:
温度 28.1
28.2
28.1
28.3
也不需要。
支付:
正常小额交易
规则评分很低
更不应该直接让 Agent 决策。
LLM 最贵的不是算力本身,而是你把高频确定性问题交给概率系统以后,连可靠性也一起变差。
我会把流水线拆成四层
Layer 0: Schema / Validity
Layer 1: Rule / Lightweight ML
Layer 2: Semantic Qualification
Layer 3: Agentic Action
Layer 0:先做最便宜的校验
例如:
JSON 是否合法
字段是否齐全
时间是否过期
是否重复事件
这层完全不要模型。
Layer 1:规则或轻量模型
例如:
金额 > 100000
状态 = FAILED
情感 = NEGATIVE
异常分 > 0.85
Google 的例子就是先用轻量 CPU 模型筛选情感,把普通事件挡掉,只把真正需要处理的负面事件送到 Agent。
Layer 2:语义资格判断
有些事件规则判断不了,但还不值得启动完整 Agent。
可以用便宜模型做:
是否需要人工调查?
是否属于已知 FAQ?
是否需要查数据库?
Layer 3:Agent
只有真正需要:
动态规划
多工具
跨系统
复杂判断
的任务才进来。
一个 Spring Boot + Kafka 的最小实现
消息模型:
public record CustomerEvent(
String eventId,
String customerId,
EventType type,
String text,
Instant occurredAt,
Map<String, Object> attributes) {
}
Pre-filter:
@Component
public class EventQualificationService {
public QualificationResult qualify(
CustomerEvent event) {
if (event.occurredAt()
.isBefore(
Instant.now()
.minus(
Duration.ofHours(24)))) {
return QualificationResult.drop(
"STALE_EVENT");
}
if (event.type()
== EventType.SYSTEM_HEARTBEAT) {
return QualificationResult.drop(
"ROUTINE_HEARTBEAT");
}
if (knownRuleMatches(event)) {
return QualificationResult.rule(
"RULE_HANDLER");
}
if (looksComplex(event)) {
return QualificationResult.agent(
"AGENT_TRIAGE");
}
return QualificationResult.normal();
}
}
Kafka Consumer:
@KafkaListener(
topics = "customer-events",
groupId = "ai-prefilter")
public void onEvent(
CustomerEvent event) {
QualificationResult result =
qualificationService
.qualify(event);
switch (result.route()) {
case DROP -> metrics.recordDrop(
result.reason());
case RULE -> ruleTopic.publish(event);
case AGENT -> agentTopic.publish(
QualifiedAgentEvent.of(
event,
result));
case NORMAL -> normalTopic.publish(event);
}
}
不要让 Pre-filter 自己变成另一个昂贵 Agent
我见过一种“优化”:
先调用一个大模型判断
需不需要再调用大模型
这经常没有意义。
Pre-filter 的目标应该是:
便宜
快
高召回
也就是说:
宁可多放一点复杂事件进去
也不要误杀真正重要事件
这是典型的 Gate 思路。
评价 Pre-filter 不能只看 Accuracy
假设:
99% 都是普通事件
一个模型永远预测“普通”:
Accuracy = 99%
但完全没用。
真正该看:
Critical Event Recall
Agent Reduction Rate
False Drop Rate
Cost Saved
Latency Added
尤其是:
False Drop Rate
高价值异常被挡在 Agent 外面,比多花一点模型钱严重得多。
一个比较实用的目标
比如告警场景:
Critical Recall >= 99.5%
Agent Traffic Reduction >= 80%
P95 Gate Latency < 20ms
False Drop Critical = 0
具体数字根据业务调整。
Pre-filter 结果也要可解释
不要只有:
0
1
记录:
public record QualificationResult(
Route route,
String reason,
double score,
String ruleVersion,
String modelVersion) {
}
以后用户问:
“这条为什么没有进入 Agent?”
才能查。
事件去重一定要在 Agent 前面
Kafka 是 At-least-once 语义时,重复消息很常见。
如果每个重复消息都启动 Agent:
成本翻倍
更危险的是重复副作用。
例如:
一条退款投诉消息重复两次
→ 两个 Agent 同时启动
→ 两边都创建补偿单
所以 Layer 0 就要:
Dedup
最简单:
insert into processed_event(
event_id,
received_at
)
values (?, now())
on conflict do nothing;
只有插入成功才继续。
Qualified Topic 不要只把原消息复制过去
最好增加 Qualification Metadata:
public record QualifiedAgentEvent(
CustomerEvent source,
String routeReason,
double priority,
RiskLevel risk,
String qualificationVersion,
Instant qualifiedAt) {
}
Agent Worker 不需要重新判断所有已知信息。
Agent Worker 再做 Budget Gate
即使通过 Pre-filter,也不代表一定马上执行。
例如:
低优先级投诉
当前模型配额耗尽
可以进入延迟队列。
if (!budgetService.canReserve(
tenantId,
estimatedCost)) {
deferQueue.publish(event);
return;
}
高吞吐里最容易被忘的是 Backpressure
Agent 一次任务可能 10 秒。
Kafka 每秒进来 500 条 Qualified Event。
如果 Worker 每秒只能处理 100:
Queue 无限增长
必须明确:
Input Rate
Qualified Rate
Agent Throughput
Queue Depth
Oldest Event Age
告警:
oldest_event_age > SLA
比只看 Queue Size 更有意义。
Queue 满了以后怎么办
不能只有:
继续堆
按任务类型定义:
DROP
DEFER
RULE_FALLBACK
HUMAN_QUEUE
SHED_LOW_PRIORITY
例如:
安全告警
→不能丢
营销情感分析
→可以延迟
低价值摘要
→可以丢弃
一个 Priority Queue
public enum AgentPriority {
CRITICAL,
HIGH,
NORMAL,
LOW
}
排序:
Priority
→ Event Time
避免低价值流量把关键 Agent 堵住。
Dataflow 的思路为什么值得迁移到非 Google 技术栈
它真正有价值的不是 Dataflow 本身,而是:
静态高吞吐计算
+
动态 Agent 处理
这两个世界应该组合,而不是替代。
传统流式系统擅长:
- 过滤;
- 聚合;
- 窗口;
- 去重;
- 状态;
- 吞吐。
Agent 擅长:
- 理解复杂语义;
- 动态规划;
- 多工具;
- 非固定流程。
把 Agent 当成每一条消息的默认 Processor,是在用最贵、最慢、最不确定的工具做最普通的事。
我会怎么压测
准备 100 万条事件:
850K routine
100K rule
40K semantic simple
10K complex
比较两套架构:
A:全部进入 Agent
B:Pre-filter + Agent
看:
总模型调用数
总Token
P50/P95
Agent Queue
Critical Recall
Tool Calls
成本
而不是只看“答案一样不一样”。
监控指标
ai_prefilter_event_total{route,reason}
ai_prefilter_latency_seconds
ai_prefilter_false_drop_total{severity}
ai_agent_qualified_rate
ai_agent_queue_depth{priority}
ai_agent_oldest_event_age_seconds{priority}
ai_agent_cost_total{route}
ai_agent_tool_call_total{tool}
我最推荐加的一个图
Raw Events
↓ 100%
Pre-filter
↓ 12%
Semantic Gate
↓ 5%
Agent
↓ 2%
Human / Side Effect
每天看这个漏斗。
如果突然变成:
Agent 20%
很可能不是业务突然复杂了,而是:
- 规则失效;
- 分类模型漂移;
- 新事件类型;
- 阈值配置错;
- 上游字段变了。
最后一个判断
Agent 最贵的优化,通常不是换一个便宜 20% 的模型。
而是:
少让 80% 根本不需要推理的事件进入 Agent。
高吞吐系统最成熟的架构一直是分层处理:确定性的事情交给确定性的系统,复杂的不确定问题才交给更昂贵的推理层。
AI 时代没有改变这个原则。
只是现在很多团队因为 LLM 太好用,暂时把它忘了。
更多推荐


所有评论(0)