昨天 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 太好用,暂时把它忘了。


Logo

AtomGit AI 社区提供模型库、数据集、Agent、Token等资源

更多推荐