ワークフロー ビルダーと実行

ワークフローは 、エグゼキューターエッジ を有向グラフに結び付け、実行を管理します。 Executor の呼び出し、メッセージ ルーティング、イベント ストリーミングを調整します。

ワークフローの構築

ワークフローは、ワークフロー構造を定義するための fluent API を提供する WorkflowBuilder クラスを使用して構築されます。

using Microsoft.Agents.AI.Workflows;

var processor = new DataProcessor();
var validator = new Validator();
var formatter = new Formatter();

// Build workflow
WorkflowBuilder builder = new(processor); // Set starting executor
builder.AddEdge(processor, validator);
builder.AddEdge(validator, formatter);
var workflow = builder.Build();

ワークフローは、 WorkflowBuilder クラスを使用して構築されます。

from agent_framework import WorkflowBuilder

processor = DataProcessor()
validator = Validator()
formatter = Formatter()

# Build workflow
builder = WorkflowBuilder(start_executor=processor)
builder.add_edge(processor, validator)
builder.add_edge(validator, formatter)
workflow = builder.build()

workflow パッケージは、エグゼキューターがエッジによって接続されるグラフベースの実行モデルを提供します。

  • Executor - 入力を受信して出力を生成する処理装置
  • Edge - 1 つの Executor の出力を別の Executor の入力に接続します
  • ビルダー - Executor とエッジを定義してワークフローを構築する
  • 実行 - 指定された入力を使用してワークフローを実行します
import (
    "github.com/microsoft/agent-framework-go/workflow"
    "github.com/microsoft/agent-framework-go/workflow/inproc"
)

uppercase := workflow.NewExecutor("UppercaseExecutor", func(input string) string {
    return strings.ToUpper(input)
}).Bind()

reverse := workflow.NewExecutor("ReverseExecutor", func(input string) string {
    runes := []rune(input)
    slices.Reverse(runes)
    return string(runes)
}).Bind()

wf, err := workflow.NewBuilder(uppercase).
    AddEdge(uppercase, reverse).
    WithOutputFrom(reverse).
    Build()
if err != nil {
    return err
}

ワークフローの実行

ワークフローでは、ストリーミング実行モードと非ストリーミング実行モードの両方がサポートされています。

using Microsoft.Agents.AI.Workflows;

// Streaming execution — get events as they happen
StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, inputMessage);
await foreach (WorkflowEvent evt in run.WatchStreamAsync())
{
    if (evt is ExecutorCompletedEvent executorComplete)
    {
        Console.WriteLine($"{executorComplete.ExecutorId}: {executorComplete.Data}");
    }

    if (evt is WorkflowOutputEvent outputEvt)
    {
        Console.WriteLine($"Workflow completed: {outputEvt.Data}");
    }
}

// Non-streaming execution — wait for completion
Run result = await InProcessExecution.RunAsync(workflow, inputMessage);
foreach (WorkflowEvent evt in result.NewEvents)
{
    if (evt is WorkflowOutputEvent outputEvt)
    {
        Console.WriteLine($"Final result: {outputEvt.Data}");
    }
}
# Streaming execution — get events as they happen
async for event in workflow.run(input_message, stream=True):
    if event.type == "output":
        print(f"Workflow completed: {event.data}")

# Non-streaming execution — wait for completion
events = await workflow.run(input_message)
print(f"Final result: {events.get_outputs()}")

イベントの発生と同時に処理したい場合は、RunStreaming を使用します。

stream, err := inproc.Default.RunStreaming(context.Background(), wf, "Hello, World!")
if err != nil {
    return err
}
defer stream.Close(context.Background())

for evt, err := range stream.WatchStream(context.Background()) {
    if err != nil {
        return err
    }
    if output, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Workflow completed: %v\n", output.Output)
    }
}

ワークフローの完了を待機し、収集されたイベントを検査する場合は、 Run を使用します。

run, err := inproc.Default.Run(context.Background(), wf, "Hello, World!")
if err != nil {
    return err
}

for evt := range run.NewEvents() {
    if output, ok := evt.(workflow.OutputEvent); ok {
        fmt.Printf("Final result: %v\n", output.Output)
    }
}

ストリーミング以外の実行によって収集された Executor イベントを調べることもできます。

for evt := range run.NewEvents() {
    if evt, ok := evt.(workflow.ExecutorCompletedEvent); ok {
        fmt.Printf("%s: %v\n", evt.ExecutorID, evt.Result)
    }
}

Tip

実行可能な完全な サンプルについては、ワークフローの例 を参照してください。

ワークフローの検証

フレームワークは、ワークフローを構築するときに包括的な検証を実行します。

  • 型の互換性: 接続された Executor 間でメッセージの種類に互換性があることを確認します
  • Graph 接続: すべての Executor が開始 Executor から到達可能であることを確認します
  • Executor バインド: すべての Executor が適切にバインドされ、インスタンス化されていることを確認します
  • エッジ検証: 重複するエッジと無効な接続のチェック

実行モデル: スーパーステップ

このフレームワークでは、変更された Pregel 実行モデル (スーパーステップ ベースの処理を含む一括同期並列 (BSP) アプローチ) を使用します。

スーパーステップのしくみ

ワークフローの実行は、個別のスーパーステップに編成されます。 各スーパーステップ:

  1. 前のスーパーステップから保留中のすべてのメッセージを収集します
  2. エッジ定義に基づいてターゲット Executor にメッセージをルーティングする
  3. スーパーステップ内ですべてのターゲット Executor を同時に実行します。
  4. すべての Executor が完了するまで待機してから進めます (同期バリア)
  5. Executor によって出力された新しいメッセージをキューに入れ、次のスーパーステップに進みます。
Superstep N:
┌─────────────────┐    ┌─────────────────┐    ┌─────────────────┐
│  Collect All    │───▶│  Route Messages │───▶│  Execute All    │
│  Pending        │    │  Based on Type  │    │  Target         │
│  Messages       │    │  & Conditions   │    │  Executors      │
└─────────────────┘    └─────────────────┘    └─────────────────┘
                                                       │
                                                       │ (barrier: wait for all)
┌─────────────────┐    ┌─────────────────┐             │
│  Start Next     │◀───│  Emit Events &  │◀────────────┘
│  Superstep      │    │  New Messages   │
└─────────────────┘    └─────────────────┘

同期バリア

最も重要な特性は、スーパーステップ間の同期バリアです。 1 つのスーパーステップ内では、トリガーされたすべての Executor が並列で実行されますが、すべての Executor が完了するまでワークフローは次のスーパーステップに進むことはありません。

これはファンアウト パターンに影響します。複数のパス (1 つは実行時間の長い Executor のチェーンを持ち、もう 1 つは実行時間の長い Executor を持つパス) にファンアウトする場合、チェーンパスは実行時間の長い Executor が完了するまで進めません。

スーパーステップの理由

BSP モデルは、次の重要な保証を提供します。

  • 確定的な実行: 同じ入力を指定すると、ワークフローは常に同じ順序で実行されます
  • 信頼性の高いチェックポイント処理: 耐障害性のために、状態はスーパーステップの境界で保存できます
  • より簡単な推論:スーパーステップ間の競合状態なし。各メッセージの一貫性のあるビューが表示されます

スーパーステップ モデルの操作

相互にブロックしない完全に独立した並列パスが必要な場合は、シーケンシャル ステップを 1 つの Executor に統合します。 step1 → step2 → step3をチェーンするのではなく、そのロジックを 1 つの Executor に結合します。 その後、両方の並列パスが 1 つのスーパーステップ内で実行されます。

次のステップ

関連トピック: