Konstruktor przepływu pracy i wdrażanie

Przepływ pracy łączy funkcje wykonawcze i krawędzie ze sobą w skierowany graf i zarządza wykonywaniem. Koordynuje wywołanie egzekutora, routing komunikatów i streaming zdarzeń.

Tworzenie przepływów pracy

Przepływy pracy są tworzone przy użyciu WorkflowBuilder klasy , która zapewnia płynny interfejs API do definiowania struktury przepływu pracy:

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

Przepływy pracy są tworzone przy użyciu WorkflowBuilder klasy :

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

Pakiet workflow udostępnia model wykonywania oparty na grafach, w którym funkcje wykonawcze są połączone za pomocą krawędzi.

  • Funkcja wykonawcza — jednostka przetwarzania, która odbiera dane wejściowe i generuje dane wyjściowe
  • Edge — łączy dane wyjściowe jednej funkcji wykonawczej z danymi wejściowymi innego
  • Builder — tworzy przepływy pracy przez definiowanie funkcji wykonawczych i krawędzi
  • Run — wykonuje przepływ pracy z podanymi danymi wejściowymi
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
}

Wykonywanie przepływu pracy

Przepływy pracy obsługują zarówno tryby przesyłania strumieniowego, jak i nieprzesyłania strumieniowego.

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

Użyj RunStreaming, gdy chcesz śledzić zdarzenia na bieżąco:

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

Użyj polecenia Run, gdy chcesz poczekać na ukończenie przepływu pracy, a następnie sprawdzić zebrane zdarzenia:

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

Można również sprawdzić zdarzenia egzekutora zebrane podczas uruchomienia bez strumieniowania:

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

Wskazówka

Zobacz przykłady przepływów pracy, aby zapoznać się z kompletnymi, gotowymi do uruchomienia przykładami.

Walidacja przepływu pracy

Struktura przeprowadza kompleksową walidację podczas tworzenia przepływów pracy:

  • Zgodność typu: zapewnia, że typy komunikatów są zgodne między połączonymi funkcjami wykonawczych
  • Łączność Graph: sprawdza, czy wszyscy wykonawcy są osiągalni od wykonawcy początkowego
  • Powiązanie wykonawców: potwierdza, że wszyscy wykonawcy są prawidłowo powiązani i zainicjalizowani
  • Walidacja krawędzi: sprawdza zduplikowane krawędzie i nieprawidłowe połączenia

Model wykonywania: superkroki

Struktura używa zmodyfikowanego modelu wykonywania Pregel — podejścia równoległego zbiorczego synchronicznego (BSP) z przetwarzaniem opartym na superkrokach.

Jak działają superkroki

Organizacja przepływu zadań odbywa się w dyskretnych superkrokach. Każdy superkrok:

  1. Zbiera wszystkie oczekujące komunikaty z poprzedniego superkroka
  2. Kieruje komunikaty do docelowych wykonawców w oparciu o definicje krawędzi
  3. Uruchamia wszystkie funkcje wykonawcze obiektów docelowych jednocześnie w ramach superkroka
  4. Czeka na ukończenie wszystkich funkcji wykonawczych przed przejściem (bariera synchronizacji)
  5. Kolejkuje wszystkie nowe komunikaty emitowane przez funkcje wykonawcze dla następnego kroku
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   │
└─────────────────┘    └─────────────────┘

Bariera synchronizacji

Najważniejszą cechą jest bariera synchronizacji między superkrokami. W ramach pojedynczego superkroka wszyscy wyzwalani wykonawcy działają równolegle, ale przepływ pracy nie przechodzi do następnego superkroku, dopóki każdy wykonawca nie zakończy swojego zadania.

Ma to wpływ na wzorce fan-out: jeśli rozgałęzisz na wiele ścieżek — jedną z łańcuchem wykonawców, a drugą z pojedynczym długotrwałym wykonawcą — ścieżka z łańcuchem wykonawców nie może przejść do momentu zakończenia działania długotrwałego wykonawcy.

Dlaczego superkroki?

Model BSP zapewnia ważne gwarancje:

  • Wykonywanie deterministyczne: biorąc pod uwagę te same dane wejściowe, przepływ pracy zawsze jest wykonywany w tej samej kolejności
  • Niezawodne tworzenie punktów kontrolnych: stan można zapisać w granicach superkroku w celu zapewnienia odporności na uszkodzenia
  • Prostsze rozumowanie: brak warunków wyścigu między superkrokami; każdy superkrok widzi spójny obraz komunikatów

Praca z modelem Superstep

Jeśli potrzebujesz naprawdę niezależnych ścieżek równoległych, które nie blokują się nawzajem, skonsoliduj sekwencyjne kroki do pojedynczego modułu wykonawczego. Zamiast łączyć łańcuch step1 → step2 → step3, połącz logikę w jedną funkcję wykonawcza. Obie ścieżki równoległe są następnie wykonywane w ramach jednego superkroka.

Następne kroki

Powiązane tematy: