ストリーミングで流すのは文字ではなく段階だ
一言でいうと
エージェントで流すべきものは、トークンだけではありません。今どのノードにいるか、何が新しく決まったか、終わったあとの状態全体が、それぞれ別のモードで出てきて、選ぶ基準が違います。
なぜ必要なのか
議事録を整理してくれるエージェントを導入したら、ユーザーから「止まったようだ」と言われました。実際にはうまく動いていました。12秒ほどかかる作業なのに、その間、画面に何もなかっただけです。
最初に浮かんだ考えは、「トークンを1文字ずつ流そう」でした。ところが、このエージェントは、モデルを最後に1回だけ呼び出します。流すトークンができる時点が、すでに終わりに近い時点です。
必要だったのは、別のものでした。「今、議事録を段落に分けています」、「決定事項を4件見つけました」のように、段階と途中の結果を流すこと。そして、それはモデルではなく、グラフが知っています。
どう動くのか
Streamingのドキュメントが定めるとおり、stream(입력, stream_mode=...)(プレースホルダーは入力です)の1行で済みます。このラボのイメージのlanggraph 0.2.60で実測した形は、次のとおりです。
| モード | イベント1つの形 | いくつ出るか |
|---|---|---|
values |
状態全体(ディクショナリ) | 入力状態1つ+スーパーステップごとに1つ |
updates |
{노드이름: 그 노드가 돌려준 갱신}(プレースホルダーはノード名と、そのノードが返した更新です) |
ノードごとに1つ |
debug |
{"type": "task"/"task_result", "step": 번호, "payload": {...}}(プレースホルダーは番号です) |
ノードごとに2つ |
custom |
ノードがwriter(...)で入れたものそのまま |
ノードが呼んだ回数だけ |
| リストで渡すと | (모드이름, 값)タプル(プレースホルダーはモード名と値です) |
選んだモードのイベントを混ぜて |
ここで混同しやすい点が2つあります。
1つ目: updatesはスーパーステップ単位ではなく、ノード単位です。並んで動くノード2つが同じスーパーステップにあっても、イベントは別々に1つずつ出ます。一方、valuesは、スーパーステップが終わったあとに1回なので、並んだノード2つの結果が統合されたあとに出ます。
2つ目: debugのstepがスーパーステップ番号です。並んで動くノードは、同じ番号を持ちます。そのため、「今何番目の段階で、その段階で何が一緒に動いているか」を知りたいときは、このモードを使います。
更新だけを集めると、状態を再構築できるか
できます。ただし、リデューサーを知っている必要があります。
for event in app.stream(입력, stream_mode="updates"):
for node, update in event.items():
for key, value in update.items():
state[key] = value # ← 이어 붙이는 열쇠에서 틀린다
traceのように連結するキーを、このように上書きすると、最後の1つだけが残ります。グラフの中ではリデューサーがしてくれていたことを、外で再構築するときは、自分で行う必要があります。
ところが、リデューサーを正しく真似ても、やはり完全には同じになりません。このラボのイメージで実測すると、次のとおりです。
그래프에 노드를 split → keypoints → actions → compose 순으로 더함
updates 로 받아 이어 붙인 trace : ['split', 'keypoints', 'actions', 'compose']
values 의 마지막 이벤트의 trace : ['split', 'actions', 'keypoints', 'compose']
このコードブロックの韓国語の部分は、ノードをsplit、keypoints、actions、composeの順にグラフへ追加したこと、1行目がupdatesで受け取って連結したtraceで、2行目がvaluesの最後のイベントのtraceであることを述べています。
並んで動く2つのノードの位置が、入れ替わっています。updatesはノードをグラフに追加した順序で出て、リデューサーは同じスーパーステップの値をノード名の順に統合するからです。どちらも毎回同じ答えを出しますが(繰り返してもぶれません)、お互いに違います。
そのため、規則はこうです。順序が意味を持つデータなら、updatesで再構築したものを最終版として使いません。画面を描くのには問題ありませんが、保存したり比較したりする最終状態は、valuesの最後のイベント(またはinvokeの答え)を使います。
その点を除けば、実務の選択は、たいてい次のように分かれます。画面が累積して描く必要があるもの(進行ログ、部分的なリスト)は、updatesで受け取ってフロントエンドがまとめ、最終状態をまるごと更新すればよいものは、valuesの最後のイベントを使います。valuesは毎回状態全体が来るので、転送量が大きくなります。状態に大きな値が入っていると、その値がスーパーステップごとにもう一度送られます。
ノードが直接出す断片
状態には残したくないけれど、人には見せたいものがあります。「3番目の段落を処理中」のような進行表示です。状態に入れると、チェックポイントに載り、再開するときにもついてきます。
customモードが、その場所です。ノードがwriter引数を受け取って、直接出します。
def compose(state, writer: StreamWriter):
writer({"stage": "compose", "points": len(state["points"])})
return {"summary": ...}
writer(...)で入れたものは、状態に残らず、聞く側にだけ届きます。そのため、進行表示・部分的なテキスト・デバッグのヒントに向いています。
流したものとinvokeの答えは同じ
valuesの最後のイベントは、invoke()が返すものと同じです。自分で突き合わせて確認でき、同じだということが重要です。画面はストリーミングで描き、記録はinvokeで別に残すように、2回動かす理由がないという意味だからです。
現場での姿
1つ目: 流すものがない時点で流そうとする。モデルを最後に1回呼び出すグラフで、トークンのストリーミングだけを付けても、前の沈黙はそのままです。段階を流す必要があります。
2つ目: valuesで大きな状態を毎ステップ出す。状態に原文が入っていると、その原文がスーパーステップごとにもう一度出ます。画面が必要とするものだけを、updatesやcustomで送るほうがよいです。
3つ目: 更新を上書きでまとめる。上の再構築の落とし穴です。足跡が1つしか残りません。
4つ目: 進行表示を状態に入れる。チェックポイントが大きくなり、再開したときに古い進行表示が復活します。
実務で本当に大切なこと
- 何を見せるかを先に決めてから、モードを選びます。モードが先ではありません。
updatesはノード単位、valuesはスーパーステップ単位、debugはスーパーステップ番号までわかります。- 更新だけで状態を再構築するには、リデューサーを知っている必要があります。
- 状態に残さないものは、
customで出します。
次のラボですること
/root/work/agstream/stream.pyを、1ステップずつ育てます。同じグラフを、values・updates・debug・複数モード同時・customの5つの方式で流しながら、イベントの数と形を自分で数えます。更新だけを集めて最後の状態を再構築してみて、連結するキーで何がずれるかを確認します。最後に、流した最後の状態とinvokeの答えが同じかを突き合わせ、何をいつ見せるかを整理して記録に残します。