Parallel:全分支执行、汇聚与可选多线程
2026年9月18日 上午11:55:10
多个专家互不依赖时,可以用一对 Parallel 节点表达“全部分支都执行,并在汇聚点等待”。
start -> dispatch
├─ tech ─┐
├─ cost ─┼─> aggregate -> decision -> end
└─ risk ─┘
1、每个分支使用独立输出键
Agent tech = createAgent("tech_analyst", "analysis.tech");
Agent cost = createAgent("cost_analyst", "analysis.cost");
Agent risk = createAgent("risk_analyst", "analysis.risk");
不要让多个分支共同写 analysis.result,否则后写入者会覆盖前面的结果。
2、构建分发和汇聚
TeamAgent reviewTeam = TeamAgent.of(null)
.name("solution_review_graph")
.protocol(TeamProtocols.NONE)
.agentAdd(tech, cost, risk)
.graphAdjuster(spec -> {
spec.addStart(Agent.ID_START).linkAdd("dispatch");
spec.addParallel("dispatch")
.linkAdd(tech.name())
.linkAdd(cost.name())
.linkAdd(risk.name());
spec.addActivity(tech).linkAdd("aggregate");
spec.addActivity(cost).linkAdd("aggregate");
spec.addActivity(risk).linkAdd("aggregate");
spec.addParallel("aggregate")
.task((ctx, node) -> {
String input = "技术:\n" + ctx.getAs("analysis.tech")
+ "\n\n成本:\n" + ctx.getAs("analysis.cost")
+ "\n\n风险:\n" + ctx.getAs("analysis.risk");
ctx.put("decision.input", input);
})
.linkAdd("decision_step");
spec.addActivity("decision_step")
.task((ctx, node) -> {
AgentSession session = ctx.getAs(Agent.KEY_SESSION);
String input = ctx.getAs("decision.input");
String result = decisionAgent.prompt(input)
.session(session)
.call()
.getContent();
ctx.put("output.final", result);
TeamTrace trace = TeamTrace.getCurrent(ctx);
if (trace != null) {
trace.setFinalAnswer(result);
}
})
.linkAdd(Agent.ID_END);
spec.addEnd(Agent.ID_END);
})
.build();
在 NONE 协议下,专家输出不会自动进入决策 Agent 的 Prompt,所以汇聚节点必须显式构造 decision.input。
3、Parallel 首先表示拓扑,不等于多线程
Parallel 保证:
- 分发节点的所有下游分支都会运行;
- 汇聚节点等待对应分支到达;
- 汇聚节点的 task 只在汇聚完成后执行。
当前 TeamAgent 内部默认执行器未配置自定义 Executor,因此上面的三个 Agent 是“并行拓扑语义”,不应承诺一定由多个线程同时请求模型。
独立 Graph 可以配置 Executor:
ExecutorService executor = Executors.newFixedThreadPool(3);
SimpleFlowDriver driver = SimpleFlowDriver.builder()
.executor(executor)
.build();
FlowEngine engine = FlowEngine.newInstance(driver);
应用关闭时应关闭自己创建的线程池。
4、Parallel 的两个关键限制
出边条件不会筛选分支
Parallel 会执行全部出边,不应在其连接上写 when(...) 并期待条件过滤。动态选择多个分支应使用 Inclusive。
汇聚必须与分发严格配对
Parallel 汇聚等待所有静态入边到达。不要让它的某条入边来自“可能不执行”的条件路径,否则汇聚可能无法继续。
推荐:
一个 Parallel split
-> 固定分支集合
一个 Parallel join
<- 同一固定分支集合
5、真正并发时的 Context 约束
默认 Context 的顶层 Map 支持并发访问,但下面的组合操作仍不安全:
List<String> results = ctx.getAs("results");
results.add(value);
分支应各写独立键,由汇聚节点统一读取。不要让多个线程共同修改普通 List,也不要共享可变请求对象。
6、不要依赖完成顺序
真正并发时,事件和记录可能交错。框架没有自动生成通用 branchId;如有需要,应以稳定 Node ID 为基础,把业务关联键写入:
branch.tech.requestId
branch.cost.requestId
branch.risk.requestId
测试只断言分支集合和汇聚结果,不断言自然完成顺序:
Assertions.assertNotNull(context.get("analysis.tech"));
Assertions.assertNotNull(context.get("analysis.cost"));
Assertions.assertNotNull(context.get("analysis.risk"));
Assertions.assertNotNull(context.get("output.final"));
关键发布动作前还应检查三份材料是否齐全,并对失败分支定义重试、降级或人工兜底。