分享一个多Agent协同中通信层设计的完整实现方案。

问题场景

某3C产线同时部署了调度Agent和质检Agent。各Agent独立运行,运行一周后问题频发:

  • 调度Agent把工单分给了已经报修的工位 → 工单积压

  • 质检Agent判定的不良品,调度Agent不知道 → 不良品继续往下流

  • 两个Agent同时决策时互相覆盖 → 信息反复横跳

产线班长反馈:“之前等通知虽然慢,但至少不会乱。现在系统快了,但乱起来更麻烦。”

问题根源

两个Agent各自独立决策,缺少一个“协同层”来同步状态信息和仲裁决策冲突。就像两个员工各自干活但不沟通,很容易互相干扰。

解决方案:基于Kafka的通信层设计

在Agent之上加一层轻量级通信中间件,负责三件事:

1. 状态广播
Agent状态变更时通知所有相关Agent。调度Agent订阅“工位状态变更”主题,收到“某工位报修”后不再往该工位分配新工单。

2. 冲突仲裁
两个Agent决策冲突时按预设优先级裁决。质检Agent的判定优先级高于调度Agent的调度指令——质量安全优先于效率。

3. 事件溯源
记录所有Agent的决策日志和通信记录。哪个Agent在什么时间做出了什么决策,全部可追溯,便于事后排查。

技术实现

// 调度Agent发布状态变更
@Service
public class DispatchAgentService {
    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;
    
    public void updateWorkstationStatus(String wsId, String status) {
        // 更新本地状态
        workstationCache.put(wsId, status);
        // 广播状态变更
        kafkaTemplate.send("topic_workstation_status", 
            new WorkstationStatusEvent(wsId, status, System.currentTimeMillis()));
    }
}

// 质检Agent订阅状态变更
@Component
public class QualityAgentListener {
    @KafkaListener(topics = "topic_workstation_status")
    public void onWorkstationStatusChange(WorkstationStatusEvent event) {
        if ("REPAIRING".equals(event.getStatus())) {
            // 不再往该工位分配新工单
            agentCache.markUnavailable(event.getWorkstationId());
            log.info("工位{}已标记为不可用,原因:报修中", event.getWorkstationId());
        }
    }
}

数据流转架构

调度Agent ──发布──→ topic_workstation_status ──订阅──→ 质检Agent
质检Agent ──发布──→ topic_inspection_result ──订阅──→ 调度Agent
调度Agent ──发布──→ topic_workstation_status ──订阅──→ 运维Agent

所有Agent不直接通信,全部通过Kafka主题交换信息,解耦同时保证数据一致性。

落地效果

  • Agent间冲突从每周7-8次降至0次

  • 信息同步延迟稳定在200ms以内

  • 事件溯源日志帮助快速定位了3次潜在的决策异常

  • 新增Agent只需订阅相关主题,无需修改现有代码

踩坑提醒

⚠️ Kafka消费者组配置要合理,避免同一Group内重复消费
⚠️ 事件消息需携带时间戳,便于事件溯源和时序判断
⚠️ 关键状态变更需配置持久化,防止Kafka宕机导致状态丢失

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐