Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Рабочий процесс связывает исполнителей и ребер между собой в направленный граф и управляет выполнением. Он координирует вызов исполнителя, маршрутизацию сообщений и потоковую передачу событий.
Создание рабочих процессов
Рабочие процессы создаются с помощью WorkflowBuilder класса, который предоставляет простой API для определения структуры рабочего процесса:
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();
Рабочие процессы создаются с помощью WorkflowBuilder класса:
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()
Пакет workflow предоставляет графовую модель выполнения, в которой исполнители соединены рёбрами.
- Исполнитель — единица обработки, которая получает входные данные и создает выходные данные
- Edge — подключает выходные данные одного исполнителя к входным данным другого
- Builder — создает рабочие процессы, задавая исполнителей и рёбра
- Выполнение — выполняет рабочий процесс с заданными входными данными
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
}
Выполнение рабочего процесса
Рабочие процессы поддерживают как режимы потоковой передачи, так и непотоковые режимы выполнения:
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()}")
Используйте RunStreaming, если хотите отслеживать события по мере их возникновения:
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)
}
}
Используйте 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)
}
}
Вы также можете просмотреть события исполнителя, собранные в ходе непотокового запуска:
for evt := range run.NewEvents() {
if evt, ok := evt.(workflow.ExecutorCompletedEvent); ok {
fmt.Printf("%s: %v\n", evt.ExecutorID, evt.Result)
}
}
Tip
См. примеры рабочих процессов, где приведены полные готовые к запуску примеры.
Проверка рабочего процесса
Платформа выполняет комплексную проверку при создании рабочих процессов:
- Совместимость типов. Гарантирует совместимость сообщений между подключенными исполнителями
- Связность графа: проверяет, достижимы ли все выполнители от начального выполнителя.
- Привязка исполнителей: подтверждает правильность привязки и начальной инициализации всех исполнителей.
- Проверка рёбер: проверяет наличие повторяющихся ребер и недопустимых подключений
Модель выполнения: supersteps
Фреймворк использует измененную модель выполнения Pregel — метод массового синхронного параллелизма (BSP) с обработкой, основанной на супершаге.
Как работают суперстепы
Выполнение рабочего процесса организовано в дискретные супершаги. Каждый супершаг:
- Собирает все ожидающие сообщения из предыдущей суперстеп
- Маршрутизирует сообщения целевым исполнителям на основе определений сетевых границ
- Выполняет все целевые исполнители одновременно в суперстепе
- Ожидает завершения всех исполнителей перед продвижением (барьер синхронизации)
- Помещает в очередь все новые сообщения, создаваемые исполнителями, для следующего суперстепа.
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 │
└─────────────────┘ └─────────────────┘
Барьер синхронизации
Наиболее важной особенностью является барьер синхронизации между суперэтапами. В рамках одного суперстепа все триггеры запускаются параллельно, но рабочий процесс не переходит к следующему суперстепу до тех пор, пока каждый исполнитель не завершит работу.
Это влияет на шаблоны разветвления: если вы разветвляете по нескольким путям — с одной стороны, последовательная цепочка выполняющих процессов, а с другой — единый длительно выполняющийся процесс — последовательный путь не может продвигаться до тех пор, пока длительно выполняющийся процесс не завершится.
Почему Суперстеп?
Модель BSP обеспечивает важные гарантии:
- Детерминированное выполнение: учитывая те же входные данные, рабочий процесс всегда выполняется в том же порядке.
- Надежный чекпоинтинг: состояние можно сохранить на границе супершага для обеспечения отказоустойчивости
- Упрощённое обоснование: отсутствие гонок между супершагами; каждый видит согласованную картину сообщений
Работа с моделью Superstep
Если вам нужны действительно независимые параллельные пути, которые не блокируют друг друга, консолидируйте последовательные шаги в один исполнитель. Вместо чейниннирования step1 → step2 → step3, объедините эту логику в одного исполнителя. Оба параллельных пути затем выполняются в одном суперстепе.
Дальнейшие действия
См. также:
- Исполнители — единицы обработки в рабочем процессе
- Ребра — связи между исполнителями
- События — наблюдаемость рабочего процесса
- Управление состоянием