並列で難しいのは分けることではなく、合わせること
一言でいうと
ブランチを増やすのは、エッジを数行書けば済みますが、統合ルールを決めないと、グラフが例外で死にます。そして、統合される順序は、偶然ではなく、決まっています。
なぜ必要なのか
1つの質問を、社内Wiki・過去のチケット・製品ドキュメントの3か所で、同時に探すようにしました。エッジを3行書けば済む作業でした。
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.
3つのノードが同じスーパーステップで同じキーに書きました。順番に動くグラフなら、最後の値が残って終わったはずですが(それも問題ですが、黙って通り過ぎます)、同時に書くと、LangGraphが「どちらを残すべきか、私にはわからない」と言って止まります。
この例外は親切なものです。順番に動くグラフで黙って上書きされていたバグが、並列に変えた瞬間に現れただけです。
どう動くのか
LangGraphの実行は、スーパーステップ単位です。1つのスーパーステップで行けるノードが複数あれば、そのノードたちが一緒に動きます。そのため、並べたノードは、いくつあっても1スーパーステップで、再帰の上限も1つしか使いません。
統合は、チャネルごとのリデューサーが行います。同じスーパーステップで、複数の値が1つのチャネルに入ってくると、
- リデューサーがあれば、その関数で1つずつまとめます。
- リデューサーがなければ、
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は、「このノードをこの入力で動かす」という指示です。通常と違う点が2つあります。
- 個数が実行時に決まります。リストの長さだけ、そのノードが動きます。
- そのノードが受け取るのは、状態全体ではなく、
Sendが渡したものです。そのため、ノードはstate["source"]のように、渡された断片だけを見ます。
戻ってくる値は、いつもどおりリデューサーでまとめられます。実測すると、結果の順序はSendを作った順序に従います。
幅を制限する
「データソースをすべて調べる」は、データソースが3つのときにだけよい考えです。20個になると、20回の呼び出しが一度に出ていき、それぞれがお金と時間です。そのため、幅(fan-out width)に上限を置きます。
上限を置くときに重要なのは、選ぶルールが決定的でなければならないことです。同点のときに適当に選ぶと、同じ質問に違う答えが出て、「昨日はできたのに今日はできない」を説明できません。スコアでソートして、同点は名前で分けるように、完全な順序を作ります。
1つが失敗すると全体が死ぬ
並列で最も痛い場所です。ブランチ1つが例外を投げると、そのスーパーステップ全体が例外で終わり、すでに成功した残りのブランチの結果も、一緒に消えます。
直し方は単純です。失敗を例外ではなく値として返します。
def node(state):
if 못 찾겠다:
return {"failures": ["ticket"]} # 예외를 던지지 않는다
return {"findings": [...]}
そうすれば、統合するノードが「3つのうち2つで見つかり、1つは失敗した」を知ることができ、ユーザーにもそう伝えられます。失敗したデータソースの名前を残すことが核心です。残さないと、部分的な結果が全体の結果のような顔をすることになります。
現場での姿
1つ目: 並列に変えたとたんに例外が出る。上のInvalidUpdateErrorです。原因は並列ではなく、もともとあった上書きです。
2つ目: 集約するノードが早く動く。ブランチごとにENDに送っておいて、集約するノードを別に置くと、集約するノードがいつ動くのかわかりません。ブランチを集約するノードに集めて送れば、そのノードは、ブランチがすべて終わったあとで、1回だけ動きます。
3つ目: 部分的な失敗が全体の失敗になる。データソース1つが死んで、答えをまったく返せないことがよく起きます。
4つ目: 幅が制御されない。リストをそのまま散らすと、リストが長くなる日に、コストがその分大きくなります。
実務で本当に大切なこと
- 並列に変える前に、「同じ欄に書くか」を先に尋ねます。リデューサーがその答えです。
- 統合される順序は決まっていますが、意味には使いません。順序が意味を持つなら、統合するノードで明示的にソートします。
- 失敗は値として返します。例外は、すでに成功したブランチまで捨てます。
- 幅に上限を置き、選ぶルールを決定的にします。
次のラボですること
/root/work/agpara/fanout.pyを、1ステップずつ育てます。まず、3つのデータソースに分かれて集まるグラフを作り、リデューサーのないキーに2つが同時に書くとどんな例外が出るかを、自分で捕まえてみます。統合される順序が何で決まるかを、名前を変えながら確認し、集約するノードを別に置いて順位を付けます。次に、Sendで個数が実行時に決まるブランチを作り、幅に上限を置き、最後に、ブランチ1つが失敗しても、残りの答えが生き残るようにします。