并行真正难的不是分开,而是合起来
一句话总结
增加分支只需要几行边,但不确定合并规则,图就会因异常而崩溃。而且合并的顺序不是偶然的,而是确定的。
为什么需要它
让一个问题同时在内部 wiki、过往工单、产品文档三个地方查找。只需要三行边的事。
for source in ("wiki", "ticket", "manual"):
graph.add_node(source, search(source))
graph.add_edge(START, source)
graph.add_edge(source, END)
第一次运行就崩了。
langgraph.errors.InvalidUpdateError: At key 'note':
Can receive only one value per step. Use an Annotated key to handle multiple values.
三个节点在同一个超级步里写入了同一个键。如果是依次运行的图,会留下最后一个值就结束了(这也是个问题,但会悄悄过去),而同时写入时,LangGraph 会说“两个里留哪个,我不知道”然后停下。
这个异常是友善的。在依次运行的图里悄悄被覆盖的 bug,只是在改成并行的那一刻暴露了出来。
工作原理
LangGraph 的执行以超级步为单位。如果一个超级步里能去的节点有多个,这些节点会一起运行。所以并排放置的节点不管有几个,都是一个超级步,也只消耗一个递归上限。
合并由各通道的 Reducer 完成。同一个超级步里有多个值进入同一个通道时,
- Reducer 有的话,用那个函数逐个合并。
- Reducer 没有的话,以
InvalidUpdateError停下。
所以改成并行之前该问的问题,不是“同时运行几个”,而是“这些分支是否写入同一格”。
合并的顺序是确定的——亲自测得的结果
这是容易误解的地方。人很容易以为“并行的话,先结束的应该排在前面”。在本实验镜像中亲自测量,并非如此。
노드를 z, a, m 순으로 더함 → 합쳐진 결과 ['a', 'm', 'z']
노드를 m, z, a 순으로 더함 → 합쳐진 결과 ['a', 'm', 'z']
'aaa' 를 일부러 느리게 만들고 'bbb' 를 즉시 끝나게 함 → ['aaa', 'bbb']
是按节点名称排序的。不是添加的顺序,也不是结束的顺序。即使慢的节点名字排在前面,它也依然排在前面。
这为什么重要?因为这意味着结果是确定的。同一输入会得到同一答案,所以可以写测试,可以比较回归。不过不要依赖名称——要让顺序承载含义,就在合并节点里显式排序。
数量在运行时才确定时——Send
数据源不固定的话,就无法提前画出边。这种时候要用 Send。
def fan_out(state):
return [Send("probe", {"source": name, "topic": state["topic"]})
for name in state["picked"]]
graph.add_conditional_edges("plan", fan_out, ["probe"])
Send 是“用这个输入运行这个节点”的指令。有两点与通常不同。
- 数量在运行时确定。列表有多长,那个节点就运行多少次。
- 那个节点接收到的不是完整状态,而是
Send传来的内容。所以节点只看到像state["source"]这样传来的片段。
返回的值照常由 Reducer 合并。亲自测量,结果顺序遵循创建 Send 的顺序。
限制宽度
“把数据源全部查一遍”只有在数据源是三个时才是个好主意。变成二十个,就会一次发出二十次调用,每一次都是钱和时间。所以要给宽度设上限(fan-out width)。
设上限时重要的是,挑选的规则必须是确定的。分数相同时随意挑选,同一个问题就会得到不同的答案,那就无法解释“昨天还行,今天就不行了”。要按分数排序、分数相同时按名称区分,这样做出完全的顺序。
一个失败,全部就崩
这是并行中最痛的地方。一个分支抛出异常,整个超级步都会以异常结束,已经成功的其余分支的结果也会一起消失。
修复的方法很简单。把失败作为值返回,而不是抛异常。
def node(state):
if 못 찾겠다:
return {"failures": ["ticket"]} # 예외를 던지지 않는다
return {"findings": [...]}
这样,合并节点就能知道“三个里有两个查到了,一个失败了”,也能这样告诉用户。关键是留下失败的数据源的名称——不留的话,部分结果就会冒充成全部结果。
在现场相遇的样子
第一,一改成并行就抛异常。就是上面的 InvalidUpdateError。原因不是并行,而是原本就存在的覆盖。
第二,汇总节点运行得太早。如果让每个分支各自流向 END,再另设一个汇总节点,就无法知道汇总节点什么时候运行。把分支们汇聚到汇总节点,那个节点就会在所有分支都结束之后只运行一次。
第三,部分失败变成整体失败。一个数据源挂了,导致根本给不出答案的事经常发生。
第四,宽度失控。把列表原样撒出去,列表变长的那天,成本就会相应增大。
实际工作中真正重要的事
- 改成并行之前,先问“是否写入同一格”。Reducer 就是答案。
- 合并的顺序是确定的,但不要用它承载含义。顺序有意义的话,就在合并节点里显式排序。
- 失败作为值返回。异常会连已经成功的分支一起丢掉。
- 给宽度设上限,并让挑选规则是确定的。
下一项实验要做什么
一步步扩展 /root/work/agpara/fanout.py。先用三个数据源做出分叉后又汇合的图,亲手捕捉没有 Reducer 的键被两个分支同时写入时会抛什么异常。通过改名字来确认合并的顺序由什么决定,另设汇总节点来排名。然后用 Send 做出数量在运行时确定的分支,给宽度设上限,最后让一个分支失败时其余答案仍然能保留下来。