
昨天 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% 都进 Agent10,000,000 次模型工作如果前面加 Gate500,000 次直接少一个数量级。而且 Agent 通常不是一次模型调用。一次复杂任务可能包含分类 → Tool → 再推理 → 再 Tool → 总结所以实际节省不是简单的 95%。最常见的错误把“智能”理解成“所有数据都送模型”例如日志告警系统CPU 42% CPU 43% CPU 41% CPU 44%这些正常数据为什么要进入 LLMIoT温度 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 ActionLayer 0先做最便宜的校验例如JSON 是否合法 字段是否齐全 时间是否过期 是否重复事件这层完全不要模型。Layer 1规则或轻量模型例如金额 100000 状态 FAILED 情感 NEGATIVE 异常分 0.85Google 的例子就是先用轻量 CPU 模型筛选情感把普通事件挡掉只把真正需要处理的负面事件送到 Agent。Layer 2语义资格判断有些事件规则判断不了但还不值得启动完整 Agent。可以用便宜模型做是否需要人工调查 是否属于已知 FAQ 是否需要查数据库Layer 3Agent只有真正需要动态规划 多工具 跨系统 复杂判断的任务才进来。一个 Spring Boot Kafka 的最小实现消息模型publicrecordCustomerEvent(StringeventId,StringcustomerId,EventTypetype,Stringtext,InstantoccurredAt,MapString,Objectattributes){}Pre-filterComponentpublicclassEventQualificationService{publicQualificationResultqualify(CustomerEventevent){if(event.occurredAt().isBefore(Instant.now().minus(Duration.ofHours(24)))){returnQualificationResult.drop(STALE_EVENT);}if(event.type()EventType.SYSTEM_HEARTBEAT){returnQualificationResult.drop(ROUTINE_HEARTBEAT);}if(knownRuleMatches(event)){returnQualificationResult.rule(RULE_HANDLER);}if(looksComplex(event)){returnQualificationResult.agent(AGENT_TRIAGE);}returnQualificationResult.normal();}}Kafka ConsumerKafkaListener(topicscustomer-events,groupIdai-prefilter)publicvoidonEvent(CustomerEventevent){QualificationResultresultqualificationService.qualify(event);switch(result.route()){caseDROP-metrics.recordDrop(result.reason());caseRULE-ruleTopic.publish(event);caseAGENT-agentTopic.publish(QualifiedAgentEvent.of(event,result));caseNORMAL-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记录publicrecordQualificationResult(Routeroute,Stringreason,doublescore,StringruleVersion,StringmodelVersion){}以后用户问“这条为什么没有进入 Agent”才能查。事件去重一定要在 Agent 前面Kafka 是 At-least-once 语义时重复消息很常见。如果每个重复消息都启动 Agent成本翻倍更危险的是重复副作用。例如一条退款投诉消息重复两次 → 两个 Agent 同时启动 → 两边都创建补偿单所以 Layer 0 就要Dedup最简单insertintoprocessed_event(event_id,received_at)values(?,now())onconflictdonothing;只有插入成功才继续。Qualified Topic 不要只把原消息复制过去最好增加 Qualification MetadatapublicrecordQualifiedAgentEvent(CustomerEventsource,StringrouteReason,doublepriority,RiskLevelrisk,StringqualificationVersion,InstantqualifiedAt){}Agent Worker 不需要重新判断所有已知信息。Agent Worker 再做 Budget Gate即使通过 Pre-filter也不代表一定马上执行。例如低优先级投诉 当前模型配额耗尽可以进入延迟队列。if(!budgetService.canReserve(tenantId,estimatedCost)){deferQueue.publish(event);return;}高吞吐里最容易被忘的是 BackpressureAgent 一次任务可能 10 秒。Kafka 每秒进来 500 条 Qualified Event。如果 Worker 每秒只能处理 100Queue 无限增长必须明确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 QueuepublicenumAgentPriority{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 BPre-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 太好用暂时把它忘了。