Arbetsflödesbyggare och utförande

Ett arbetsflöde kopplar exekutorer och kanter till en riktad graf och hanterar flödet. Den samordnar anrop av exekutorer, meddelanderoutning och händelseströmning.

Skapa arbetsflöden

Arbetsflöden skapas med hjälp av WorkflowBuilder klassen, vilket ger ett flytande API för att definiera arbetsflödesstrukturen:

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

Arbetsflöden skapas med hjälp av WorkflowBuilder klassen:

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

Paketet workflow innehåller en grafbaserad körningsmodell där utförare är anslutna via kanter.

  • Executor – en bearbetningsenhet som tar emot indata och genererar utdata
  • Edge – Ansluter utdata från en köre till indata från en annan
  • Builder – Bygger arbetsflöden genom att definiera exekverare och kanter
  • Kör – kör ett arbetsflöde med angivna indata
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
}

Arbetsflödeskörning

Arbetsflöden stöder både körningslägen för direktuppspelning och icke-direktuppspelning:

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

Använd RunStreaming när du vill ha händelser när de inträffar:

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

Använd Run när du vill vänta tills arbetsflödet har slutförts och kontrollera sedan de insamlade händelserna:

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

Du kan också granska händelser från exekveraren som har samlats in under en körning utan strömning:

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

Tips/Råd

Se arbetsflödesexemplen för fullständiga körbara exempel.

Arbetsflödesverifiering

Ramverket utför omfattande validering när arbetsflöden skapas:

  • Typkompatibilitet: Säkerställer att meddelandetyper är kompatibla mellan anslutna körverktyg
  • Graph Connectivity: Verifierar att alla exekveringskomponenter är nåbara från startkomponenten.
  • Exekutorbindning: Bekräftar att alla exekutorer är korrekt bundna och instansierade
  • Kantverifiering: Söker efter duplicerade kanter och ogiltiga anslutningar

Körningsmodell: Supersteg

Ramverket använder en modifierad Pregel-körningsmodell – en Bulk Synchronous Parallel (BSP)-ansats med bearbetning baserad på supersteg.

Så här fungerar Supersteps

Arbetsflödeskörning ordnas i diskreta supersteg. Varje supersteg:

  1. Samlar in alla väntande meddelanden från föregående supersteg
  2. Dirigerar meddelanden till målexekutorer baserat på gränsdefinitioner
  3. Kör alla målexekutorer samtidigt i supersteget
  4. Väntar på att alla exekutörer har slutfört sitt arbete innan de avancerar (synkroniseringsbarriär)
  5. Lägger i kö alla nya meddelanden som genereras av exekutorer för nästa supersteg
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   │
└─────────────────┘    └─────────────────┘

Synkroniseringsbarriär

Den viktigaste egenskapen är synkroniseringsbarriären mellan supersteg. Inom ett enda supersteg körs alla aktiverade utförare parallellt, men arbetsflödet går inte vidare till nästa supersteg förrän varje utförare har slutfört sitt uppdrag.

Detta påverkar förgreningsmönster: om du delar upp till flera sökvägar – en med en kedja av exekutorer och en annan med en enda långvarig exekutor – kan den kedjade sökvägen inte avancera förrän den långvariga exekutorn har slutförts.

Varför Supersteps?

BSP-modellen ger viktiga garantier:

  • Deterministisk körning: Med samma indata körs arbetsflödet alltid i samma ordning
  • Tillförlitlig kontrollpunkt: Systemtillstånd kan sparas vid superstegsgränser för feltolerans
  • Enklare resonemang: Inga konkurrensförhållanden mellan supersteg; var och en ser en konsekvent vy över meddelanden

Arbeta med Superstep-modellen

Om du behöver verkligt oberoende parallella sökvägar som inte blockerar varandra, konsoliderar du sekventiella steg till en enda utförare. I stället för att länka step1 → step2 → step3 bör du kombinera den logiken i en enskild exekverare. Båda parallella vägarna körs sedan inom ett enda supersteg.

Nästa steg

Relaterade ämnen: