İş Akışı Oluşturucusu ve Yürütme

İş Akışı, yürütücüleri ve kenarları yönlendirilmiş bir grafik içinde birleştirir ve yürütmeyi yönetir. Yürütücü çağırmayı, ileti yönlendirmeyi ve olay akışını koordine eder.

İş Akışları Oluşturma

İş akışları, iş akışı yapısını tanımlamak için akıcı bir API sağlayan sınıfı kullanılarak WorkflowBuilder oluşturulur:

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

İş akışları sınıfı kullanılarak WorkflowBuilder oluşturulur:

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, workflow yürütücülerin kenarlarla bağlandığı graf tabanlı bir yürütme modeli sağlar.

  • Yürütücü - Giriş alan ve çıkış üreten bir işleme birimi
  • Edge - Bir yürütücü çıktısını başka bir yürütücü girişine bağlar
  • Oluşturucu - Yürütücüleri ve kenarları tanımlayarak iş akışları oluşturur
  • Çalıştır - Verilen girişle bir iş akışı yürütür
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
}

İş Akışı Yürütme

İş akışları hem akış hem de akış dışı yürütme modlarını destekler:

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

Olayları gerçekleşirken görmek istediğinizde RunStreaming kullanın:

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

İş akışının tamamlanmasını beklemek ve ardından toplanan olayları incelemek istediğinizde kullanın 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)
    }
}

Akış dışı bir çalıştırma tarafından toplanan yürütücü olaylarını da inceleyebilirsiniz:

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

İpucu

Çalıştırılabilir örneklerin tamamı için iş akışı örneklerine bakın.

İş Akışı Doğrulama

Çerçeve, iş akışları oluştururken kapsamlı doğrulama gerçekleştirir:

  • Tür Uyumluluğu: İleti türlerinin bağlı yürütücüler arasında uyumlu olmasını sağlar
  • Graf Bağlantısı: Tüm yürütücülerin başlangıç yürütücüsundan erişilebilir olduğunu doğrular
  • Yürütücü Bağlaması: Tüm yürütücülerin düzgün şekilde bağlı ve örneklendirildiğini onaylar
  • Edge Doğrulaması: Yinelenen kenarları ve geçersiz bağlantıları denetler

Yürütme Modeli: Supersteps

Çerçeve, üst düzey tabanlı işlemeye sahip Toplu Zaman Uyumlu Paralel (BSP) yaklaşımı olan değiştirilmiş bir Pregel yürütme modeli kullanır.

Supersteps Nasıl Çalışır?

İş akışının yürütülmesi ayrı süper adımlar olarak düzenlenir. Her üst adım:

  1. Önceki üst adımdan tüm bekleyen iletileri toplar
  2. uç tanımlarına göre iletileri hedef yürütücülere yönlendirir
  3. Üst adım içinde tüm hedef yürütücüleri eşzamanlı olarak çalıştırır
  4. Tüm yürütücülerin ilerlemeden önce tamamlanmasını bekler (senkronizasyon bariyeri)
  5. Yürütücüler tarafından gönderilen tüm yeni iletileri sonraki üst adım için kuyruğa alır
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   │
└─────────────────┘    └─────────────────┘

Eşitleme Engeli

En önemli özellik, süper adımlar arasındaki senkronizasyon engelidir. Tek bir üst adım içinde, tetiklenen tüm yürütücüler paralel olarak çalışır, ancak her yürütücü tamamlanana kadar iş akışı bir sonraki üst adıma ilerlemez.

Bu, çok yönlü desenleri etkiler: biri yürütücü zincirine, diğeri ise uzun süre çalışan tek bir yürütücüye sahip birden çok yola doğru ilerlerseniz, zincirlenmiş yol uzun süre çalışan yürütücü tamamlanana kadar ilerletilemez.

Neden Supersteps?

BSP modeli önemli garantiler sağlar:

  • Deterministik yürütme: Aynı giriş göz önüne alındığında iş akışı her zaman aynı sırada yürütülür
  • Güvenilir denetim noktası oluşturma: Durum, hataya dayanıklılık için süper adım sınırlarında kaydedilebilir
  • Daha basit mantık: Süpersteps arasında yarış koşulu yoktur; her biri iletilerin tutarlı bir görünümünü görür

Superstep Modeli ile çalışma

Birbirini engellemeyen gerçekten bağımsız paralel yollara ihtiyacınız varsa, sıralı adımları tek bir yürütücüde birleştirin. zincirleme step1 → step2 → step3 yapmak yerine, bu mantığı tek bir yürütücüde birleştirin. Her iki paralel yol da tek bir üst adım içinde yürütülür.

Sonraki Adımlar

İlgili konular: