Pembuat & Eksekusi Alur Kerja

Alur Kerja mengikat eksekutor dan tepi bersama-sama ke dalam grafik yang diarahkan dan mengelola eksekusi. Ini mengoordinasikan pemanggilan pelaksana, perutean pesan, dan streaming peristiwa.

Membangun Alur Kerja

Alur kerja dibangun menggunakan WorkflowBuilder kelas , yang menyediakan API yang fasih untuk menentukan struktur alur kerja:

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

Alur kerja dibangun menggunakan WorkflowBuilder kelas :

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

Paket ini workflow menyediakan model eksekusi berbasis grafik di mana pelaksana dihubungkan oleh tepi.

  • Pelaksana - Unit pemrosesan yang menerima input dan menghasilkan output
  • Edge - Menyambungkan output dari satu eksekutor ke input pelaksana lain
  • Builder - Membangun alur kerja dengan menentukan pelaksana dan tepi
  • Jalankan - Menjalankan alur kerja dengan input yang diberikan
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
}

Eksekusi Alur Kerja

Alur kerja mendukung mode eksekusi streaming dan non-streaming:

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

Gunakan RunStreaming saat Anda ingin mengetahui kejadian secara langsung:

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

Gunakan Run saat Anda ingin menunggu penyelesaian alur kerja lalu periksa peristiwa yang dikumpulkan:

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

Anda juga dapat meninjau peristiwa eksekutor yang dikumpulkan selama eksekusi non-streaming:

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

Tip

Lihat contoh alur kerja untuk sampel lengkap yang dapat dijalankan.

Validasi Alur Kerja

Kerangka kerja melakukan validasi komprehensif saat membangun alur kerja:

  • Kompatibilitas Jenis: Memastikan jenis pesan kompatibel antara eksekutor yang terhubung
  • Konektivitas Graf: Memverifikasi semua eksekutor dapat dijangkau dari eksekutor awal
  • Pengikatan Eksekutor: Mengonfirmasi bahwa semua eksekutor terikat dan diinstansiasi dengan benar
  • Validasi Edge: Memeriksa tepi duplikat dan koneksi yang tidak valid

Model Eksekusi: Supersteps

Kerangka kerja menggunakan model eksekusi Pregel yang dimodifikasi - pendekatan Paralel Sinkron Massal (BSP) dengan pemrosesan berbasis superstep.

Cara Kerja Supersteps

Eksekusi alur kerja diatur menjadi "supersteps" diskrit. Setiap langkah super:

  1. Mengumpulkan semua pesan yang masih tertunda dari superstep sebelumnya
  2. Mengalihkan pesan ke pelaksana target berdasarkan definisi batas
  3. Menjalankan semua eksekutor target secara bersamaan dalam superstep
  4. Menunggu semua pelaksana selesai sebelum maju (hambatan sinkronisasi)
  5. Mengantre pesan baru apa pun yang dihasilkan oleh eksekutor untuk superstep berikutnya
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   │
└─────────────────┘    └─────────────────┘

Hambatan Sinkronisasi

Karakteristik yang paling penting adalah hambatan sinkronisasi antara supersteps. Dalam satu superstep, semua eksekutor yang dipicu berjalan secara paralel, tetapi alur kerja tidak maju ke superstep berikutnya sampai setiap pelaksana selesai.

Ini memengaruhi pola fan-out: jika Anda membagi ke beberapa jalur — satu dengan rangkaian eksekutor dan satu lagi dengan eksekutor tunggal yang berjalan lama — jalur dengan rantai tidak dapat maju sampai eksekutor jangka panjang selesai.

Mengapa Supersteps?

Model BSP memberikan jaminan penting:

  • Eksekusi deterministik: Mengingat input yang sama, alur kerja selalu dijalankan dalam urutan yang sama
  • Titik pemeriksaan yang andal: Status dapat disimpan pada batas superstep untuk toleransi kesalahan
  • Penalaran yang lebih sederhana: Tidak ada kondisi balapan antara superstep; setiap superstep melihat tampilan pesan yang konsisten

Bekerja dengan Model Superstep

Jika Anda memerlukan jalur paralel yang benar-benar independen yang tidak saling memblokir, konsolidasikan langkah berurutan ke dalam satu pelaksana. Alih-alih menautkan, gabungkan logika tersebut menjadi satu eksekutor. Kedua jalur paralel kemudian dijalankan dalam satu superstep.

Langkah berikutnya

Topik terkait: