别把每条 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 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  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 太好用,暂时把它忘了。


更多企业级 AI 应用、Agent、RAG 与模型工程化内容,我会继续整理在 智元界

https://www.zyentor.com/