并行执行:Fan-out / Fan-in 与并发状态合并

本文是「LangGraph 核心与实战」系列专栏的第 5 篇。专栏总览参见:《LangGraph 核心架构全景:用状态图构建可控 Agent 工作流》。

在顺序图与基础条件分支中,节点的执行在时序上是完全串行的:步骤 A 跑完才启动步骤 B。然而在许多现实业务中,任务包含多个彼此没有数据依赖的独立子环节:

  • 分析某份技术方案时,需要同时进行技术可行性评估、合规安全审查与云资源成本估算;
  • 调研竞争对手时,需要并发检索 3 个不同的数据源并分别提炼指标;
  • 多模态处理时,音频转写、图片 OCR 与文本实体抽取可以同时启动。

如果依然使用串行流水线,端到端延迟将是各个子任务耗时的累加( )。通过构建 Fan-out(扇出并发)与 Fan-in(扇入聚合)拓扑,任务耗时可以大幅降低至最慢单节点的耗时( )。


拓扑设计:Fan-out / Fan-in 模式

在状态图中表达并行执行非常直观:从同一个源节点同时引出多条普通边指向不同目标,并在后续让它们共同汇聚到一个聚合节点:

在代码声明层面,只需要添加对应的边:

1
2
3
4
5
6
7
8
9
# 1. 扇出声明 (Fan-out)
builder.add_edge("prepare", "eval_arch")
builder.add_edge("prepare", "eval_security")
builder.add_edge("prepare", "eval_cost")

# 2. 扇入声明 (Fan-in)
builder.add_edge("eval_arch", "aggregate")
builder.add_edge("eval_security", "aggregate")
builder.add_edge("eval_cost", "aggregate")

LangGraph 运行时提供了关键的执行保证:汇聚节点 aggregate 只有在所有上游并发节点(eval_arch、eval_security、eval_cost)全部执行完成并成功更新状态后,才会被调度触发。


核心并发陷阱:写冲突(Race Condition)与 Reducer

在并行执行中,必须警惕状态更新的竞态冲突。

为什么默认字典合并会丢失数据

回顾前两篇提到的状态合并规则:LangGraph 默认采用浅层字段覆写。
假设状态中定义了普通字典或普通列表字段:

1
2
class UnsafeState(TypedDict):
evaluations: list # 未声明 Reducer 的普通列表

当 eval_arch 与 eval_security 并发执行时:

  • eval_arch 完成,向状态提交 {"evaluations": [{"type": "arch", "score": 90}]};
  • eval_security 几乎在同一时刻完成,向状态提交 {"evaluations": [{"type": "security", "score": 85}]};
  • 由于两个节点在各自的执行副本中拿到的是旧 State,且字段策略是直接覆盖,后完成的节点会直接覆盖先完成节点的数据,导致先完成的评估结果彻底蒸发!

解决方案:使用 Annotated + operator.add

为了让并发节点能够安全地向同一容器追加数据,必须通过 operator.add 为该字段声明归约合并器(Reducer):

1
2
3
4
5
6
7
8
from typing import Annotated, List, TypedDict, Dict, Any
import operator

class SafeState(TypedDict):
target_name: str
# 核心:必须声明 Reducer,确保并发写入时按列表合并而不是相互覆盖
eval_reports: Annotated[List[Dict[str, Any]], operator.add]
final_decision: str

有了 Annotated[..., operator.add],LangGraph 在运行时会接收所有并发节点的增量返回值,按内部线程同步机制安全执行追加操作,汇聚节点 aggregate 最终看到的将是包含所有并发结果的完整列表。


实战:三维度方案综合评审工作流

实现一个包含准备、并发三路审查与终审决策的完整可运行示例:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
import operator
from typing import TypedDict, Annotated, List, Dict, Any
from langgraph.graph import StateGraph, START, END

# 1. 状态契约定义
class ProposalReviewState(TypedDict):
proposal_title: str
proposal_content: str
# 并发收集的各维度评审切片
dimension_reports: Annotated[List[Dict[str, Any]], operator.add]
# 汇总产物
total_score: float
approval_status: str

# 2. 准备阶段节点
def prepare_proposal_node(state: ProposalReviewState) -> dict:
print(f"\n[工位 1] 收到方案: 《{state['proposal_title']}》,准备派发评审任务...")
return {}

# 3. 三个平行的评审节点 (并发执行)
def architecture_review_node(state: ProposalReviewState) -> dict:
content = state["proposal_content"]
print(" -> [并发审查] 架构评审节点启动...")
# 模拟架构评分
score = 90.0 if "高可用" in content else 70.0
return {
"dimension_reports": [{
"dimension": "架构设计",
"score": score,
"comment": "满足分布式容灾标准" if score >= 80 else "缺少容灾备份机制"
}]
}

def security_review_node(state: ProposalReviewState) -> dict:
content = state["proposal_content"]
print(" -> [并发审查] 安全合规节点启动...")
# 模拟安全评分
score = 85.0 if "鉴权" in content else 60.0
return {
"dimension_reports": [{
"dimension": "安全合规",
"score": score,
"comment": "数据敏感字段已配置加密" if score >= 80 else "未发现权限校验设计"
}]
}

def cost_review_node(state: ProposalReviewState) -> dict:
content = state["proposal_content"]
print(" -> [并发审查] 成本核算节点启动...")
# 模拟成本评分
score = 80.0
return {
"dimension_reports": [{
"dimension": "资源成本",
"score": score,
"comment": "云资源预算在年度配额限制内"
}]
}

# 4. 汇聚节点 (Fan-in 聚合)
def aggregate_decision_node(state: ProposalReviewState) -> dict:
reports = state["dimension_reports"]
print(f"\n[工位 3] 收到全部 {len(reports)} 个维度的并发报告,开始裁决...")

avg_score = sum(r["score"] for r in reports) / len(reports)
status = "APPROVED" if avg_score >= 80.0 else "REJECTED"

return {
"total_score": round(avg_score, 1),
"approval_status": status
}

# 5. 编排图拓扑
builder = StateGraph(ProposalReviewState)

builder.add_node("prepare", prepare_proposal_node)
builder.add_node("eval_arch", architecture_review_node)
builder.add_node("eval_security", security_review_node)
builder.add_node("eval_cost", cost_review_node)
builder.add_node("aggregate", aggregate_decision_node)

builder.add_edge(START, "prepare")

# 扇出: prepare 同时触发三个节点
builder.add_edge("prepare", "eval_arch")
builder.add_edge("prepare", "eval_security")
builder.add_edge("prepare", "eval_cost")

# 扇入: 三个节点全部连接至 aggregate
builder.add_edge("eval_arch", "aggregate")
builder.add_edge("eval_security", "aggregate")
builder.add_edge("eval_cost", "aggregate")

builder.add_edge("aggregate", END)

app = builder.compile()

# 6. 执行测试
if __name__ == "__main__":
proposal = {
"proposal_title": "核心交易链路改造方案",
"proposal_content": "本系统基于云原生多可用区部署,具备高可用容灾能力,并配置细粒度接口鉴权与日志审计。"
}

result = app.invoke(proposal)

print("\n================ 评审交付报告 ================")
print(f"方案: {result['proposal_title']}")
print(f"综合评分: {result['total_score']} | 最终决策: {result['approval_status']}")
print("各维度明细:")
for r in result["dimension_reports"]:
print(f" * [{r['dimension']}] 得分: {r['score']} -> 评价: {r['comment']}")
print("==============================================")

输出日志:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
[工位 1] 收到方案: 《核心交易链路改造方案》,准备派发评审任务...
-> [并发审查] 架构评审节点启动...
-> [并发审查] 安全合规节点启动...
-> [并发审查] 成本核算节点启动...

[工位 3] 收到全部 3 个维度的并发报告,开始裁决...

================ 评审交付报告 ================
方案: 核心交易链路改造方案
综合评分: 85.0 | 最终决策: APPROVED
各维度明细:
* [架构设计] 得分: 90.0 -> 评价: 满足分布式容灾标准
* [安全合规] 得分: 85.0 -> 评价: 数据敏感字段已配置加密
* [资源成本] 得分: 80.0 -> 评价: 云资源预算在年度配额限制内
==============================================

生产级约束:并行执行的注意事项

在生产环境中编排高并发工作流时,需遵守以下三条工程准则:

  1. 防止下游大模型 API 触发限流(Rate Limit 429):如果在一个 Fan-out 拓扑中同时扇出 20 个调用大模型的子节点,瞬时并发很可能直接打满服务商的 TPM / RPM 配额。建议结合信号量(asyncio.Semaphore)限制并发上限,或使用队列解耦;
  2. 节点间严禁隐式数据依赖:位于同一扇出层级的并发节点之间在运行时是互不可见的。如果节点 B 的计算依赖节点 A 的输出,两者必须设计为串行先后关系,绝不能放在同一个并行层级;
  3. 复合拓扑应用:Fan-in 汇聚节点产出数据后,完全可以无缝衔接上一篇所讲的条件边(例如:若综合评分低于 80 分,则通过条件边路由给退回修改节点,否则直接交付)。

系列导航与参考