Tvůrce pracovních postupů a provádění

Pracovní postup spojuje exekutory a hrany do řízeného grafu a spravuje provádění. Koordinuje vyvolání exekutoru, směrování zpráv a streamování událostí.

Vytváření pracovních postupů

Pracovní postupy se vytvářejí pomocí WorkflowBuilder třídy, která poskytuje rozhraní API fluent pro definování struktury pracovního postupu:

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();

Pracovní postupy se vytvářejí pomocí WorkflowBuilder třídy:

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()

Balíček workflow poskytuje grafový model provádění, ve kterém jsou vykonavatelé propojeni pomocí hran.

  • Exekutor – výpočetní jednotka, která přijímá vstup a vytváří výstup
  • Edge – připojí výstup jednoho exekutoru ke vstupu jiného exekutoru.
  • Tvůrce – Vytváření pracovních postupů definováním exekutorů a hran
  • Spustit - Spustí workflow se zadaným vstupem
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
}

Provádění pracovního postupu

Pracovní postupy podporují režimy streamování i spouštění bez streamování:

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()}")

Použijte RunStreaming , když chcete, aby události probíhaly takto:

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)
    }
}

Použijte Run , když chcete počkat na dokončení pracovního postupu a pak zkontrolovat shromážděné události:

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)
    }
}

Můžete také zkontrolovat události exekutoru shromážděné spuštěním bez streamování:

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

Tip

Kompletní spustitelné ukázky najdete v příkladech pracovního postupu .

Ověření pracovního postupu

Architektura provádí komplexní ověřování při vytváření pracovních postupů:

  • Kompatibilita typů: Zajišťuje kompatibilitu typů zpráv mezi připojenými exekutory.
  • Spojitost grafu: Ověřuje, že všechny exekutory jsou dostupné od spouštěcího exekutoru.
  • Vazba exekutoru: Potvrdí, že jsou všechny exekutory správně svázané a instancovány.
  • Ověření Edge: Detekuje duplicitní hrany a neplatná připojení.

Model provádění: Supersteps

Architektura používá upravený model provádění Pregel – metodu BSP (Bulk Synchronous Parallel) se zpracováním založeným na superkrocích.

Jak fungují superkroky

Provádění pracovního postupu je uspořádané do samostatných superkroků. Každý superkrok:

  1. Shromažďuje všechny čekající zprávy z předchozího superkroku.
  2. Směruje zprávy do cílových vykonavatelů na základě definic hran.
  3. Spustí všechny cílové vykonavatele souběžně v rámci superkroku.
  4. Čeká na dokončení všech vykonavatelů než postoupí (synchronizační bariéra).
  5. Zařadí všechny nové zprávy generované exekutory do fronty pro další superkrok.
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   │
└─────────────────┘    └─────────────────┘

Synchronizační bariéra

Nejdůležitější vlastností je synchronizační bariéra mezi superkroky. V rámci jednoho superkroku se všechny aktivované exekutory spouští paralelně, ale pracovní postup nepřesáhne k dalšímu superkroku, dokud se všechny exekutory nedokončí.

To má vliv na vzory rozvětvení: pokud rozvětvíte do více cest – jednu s řetězem vykonavatelů a druhou s jedním dlouhotrvajícím vykonavatelem – zřetězená cesta nemůže pokračovat, dokud dlouhotrvající vykonavatel nedokončí svou činnost.

Proč superkroky?

Model BSP poskytuje důležité záruky:

  • Deterministické provádění: Vzhledem ke stejnému vstupu se pracovní postup vždy provede ve stejném pořadí.
  • Spolehlivé vytváření kontrolních bodů: Stav lze uložit na hranicích superkroku pro odolnost proti chybám.
  • Jednodušší odůvodnění: Žádné podmínky závodu mezi superkroky; každý vidí konzistentní zobrazení zpráv.

Práce s modelem Superstep

Pokud potřebujete skutečně nezávislé paralelní cesty, které se navzájem nezablokují, konsolidujte sekvenční kroky do jednoho exekutoru. Místo řetězení step1 → step2 → step3zkombinujte tuto logiku do jednoho exekutoru. Obě paralelní cesty se pak provádějí v rámci jednoho superkroku.

Další kroky

Související témata: