Режимы выполнения рабочего процесса

При выполнении рабочего процесса в .NET режим выполнения определяет процесс обработки надстроек и способ доставки событий потребителю. Класс InProcessExecution предоставляет два режима выполнения: OffThread и Lockstep.

Overview

OffThread (по умолчанию) Блокировка
Выполнение суперстепа Фоновый поток Клиентский поток
Доставка событий События обрабатываются немедленно при их возникновении Пакетная обработка после завершения каждого суперстепа
Выполнение шага Независимо от обработки событий Приостановлено, пока пакетные события не будут обработаны
Параллелизм Потребитель считывает события во время выполнения супершагов Чередующееся выполнение потребителей и супершага
лучше всего подходит для Потоковая передача в режиме реального времени, рабочие сценарии Тестирование, отладка, детерминированное упорядочение

OffThread

OffThread — это режим выполнения по умолчанию . Суперэтапы выполняются в фоновом потоке, а события стримятся немедленно, как только они генерируются через канал основанной реализации.

// OffThread is the default — these are equivalent:
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input);
await using StreamingRun run = await InProcessExecution.OffThread.RunStreamingAsync(workflow, input);

Принцип работы

  1. Фоновая задача выполняет супершаги непрерывно, пока сообщения находятся в ожидании.
  2. При генерации выходных данных или событий результирующие WorkflowEvent объекты записываются в неограниченный Channel<WorkflowEvent>.
  3. Потребитель считывает события из канала через WatchStreamAsync, получая их в режиме реального времени по мере их производства.
  4. Когда все суперэтапы завершены и не остаётся сообщений, выполнение останавливается со статусом Idle или PendingRequests.

Так как цикл суперстепа и потребитель выполняются одновременно, события становятся видимыми сразу после их создания — буферизация не возникает. Это делает OffThread идеальным для сценариев потоковой передачи, где важна низкая задержка в доставке событий, например, для поэтапного обновления токенов в пользовательском интерфейсе.

Одновременные запуски

OffThread также поддерживает одновременный вариант, позволяющий нескольким запускам совместно использовать один экземпляр рабочего процесса одновременно:

await using StreamingRun run = await InProcessExecution.Concurrent.RunStreamingAsync(workflow, input);

Important

Для одновременного выполнения требуется, чтобы все исполнители в рабочем процессе объявлялись crossRunShareable (на конструкторе) или предоставлялись в качестве фабричных методов.

Блокировка

В режиме Lockstep супершаги выполняются в потоке потребителя, а не в фоновой задаче. События накапливаются во время каждого супершагa и испускаются в виде пакета после завершения супершагa.

await using StreamingRun run = await InProcessExecution.Lockstep.RunStreamingAsync(workflow, input);

Принцип работы

  1. Потребитель вызывает WatchStreamAsync, который управляет циклом выполнения.
  2. Супершаг выполняется до завершения, и события накапливаются в очереди.
  3. После завершения этапа супершага все события из очереди передаются потребителю.
  4. Следующий супершаг начинается только после того, как потребитель получил все события от предыдущего супершага.

Этот чередующийся шаблон означает, что потребитель и движок рабочего процесса никогда не выполняются одновременно. Доставка событий является детерминированной — все события из супершага гарантированно поступают до любых событий из следующего супершага.

Когда следует использовать Lockstep

Локстеп полезен в следующих ситуациях:

  • Тестирование — детерминированное упорядочивание событий упрощает утверждения.
  • Отладка — пошаговая отладка упрощается, когда выполнение остается в потоке потребителя.
  • Упорядоченная обработка — сценарии, в которых необходимо полностью обработать события суперстепа перед началом следующей суперстеп.

Выбор режима выполнения

Для большинства рабочих сценариев рекомендуется использовать режим offThread по умолчанию. Он обеспечивает лучшую скорость реагирования и позволяет рабочему процессу продолжать обработку, пока потребитель обрабатывает события.

Используйте lockstep , если детерминированное поведение более важно, чем производительность, например в модульных тестах или сеансах отладки.

// Production: OffThread (default)
await using StreamingRun run = await InProcessExecution.RunStreamingAsync(workflow, input);

// Testing: Lockstep for deterministic behavior
await using StreamingRun run = await InProcessExecution.Lockstep.RunStreamingAsync(workflow, input);

Операция без потокового режима

Оба режима выполнения поддерживают выполнение в нестриминговом режиме с использованием RunAsync. В режиме без потоковой передачи рабочий процесс выполняется до завершения и собирает все события в Run объект, а не выполняет их добавочно:

Run run = await InProcessExecution.RunAsync(workflow, input);

// Access all emitted events
foreach (WorkflowEvent evt in run.OutgoingEvents)
{
    // Process events
}

Так как непотоковое выполнение собирает все события после завершения, преимущество OffThread в доставке событий в режиме реального времени не применяется. Основное различие между режимами в сценариях без потоковой передачи — многопоточность: OffThread выполняет супершаги на фоновом потоке, освобождая вызывающий поток во время ожидания завершения, в то время как Lockstep выполняет супершаги на потоке вызывающего объекта, блокируя его поток до завершения рабочего процесса.

Выполнение без потоковой передачи использует режим OffThread по умолчанию. Чтобы использовать Lockstep с непотоковым выполнением, выполните следующие действия:

Run run = await InProcessExecution.Lockstep.RunAsync(workflow, input);

Дальнейшие действия

Режимы выполнения не применимы к рабочим процессам Python. Рабочие процессы Python используют единую модель выполнения, которая обрабатывает супершаги и доставляет события через асинхронный генератор. Эта модель аналогична режиму Lockstep в .NET — шаги не продвигаются, пока потребитель активно извлекает события из генератора.

С более подробной информацией о выполнении рабочих процессов на Python можно ознакомиться в Конструкторе и выполнении рабочих процессов.

При запуске рабочего процесса в Go среда выполнения определяет, как обрабатываются суперэтапы и как события доставляются потребителю. Пакет workflow/inproc предоставляет три среды: Default/OffThread, Lockstepи .Concurrent

Overview

OffThread / По умолчанию Блокировка Одновременный
Выполнение суперстепа Фоновая горутина Управляется потребителем событий Фоновая горутина
Доставка событий События обрабатываются немедленно при их возникновении Пакетная обработка по мере использования потока События обрабатываются немедленно при их возникновении
лучше всего подходит для Потоковая передача в режиме реального времени, рабочие сценарии Тестирование, отладка, детерминированное упорядочение Общие экземпляры рабочих процессов с одновременными безопасными привязками

OffThread

OffThread — это режим выполнения по умолчанию. Это эквивалентно:

stream, err := inproc.Default.RunStreaming(ctx, wf, input)
stream, err := inproc.OffThread.RunStreaming(ctx, wf, input)

Принцип работы

  1. Фоновая горутина выполняет супершаги, пока есть ожидающие обработки сообщения.
  2. По мере того как исполнители получают выходные данные или события, события рабочего процесса записываются в поток.
  3. Потребитель читает события из WatchStream, получая их по мере появления.
  4. Когда все суперэтапы завершены и не остаётся сообщений, выполнение завершается со статусом простоя или ожидания запроса.

Одновременные запуски

Используйте inproc.Concurrent , когда все привязки исполнителя в рабочем процессе поддерживают параллельное общее выполнение:

stream, err := inproc.Concurrent.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer stream.Close(ctx)

Блокировка

В режиме Lockstep выполнение рабочего процесса продвигается по мере того, как потребитель считывает данные из потока. Это делает упорядочение событий детерминированным для тестов и отладки.

stream, err := inproc.Lockstep.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer stream.Close(ctx)

for evt, err := range stream.WatchStream(ctx) {
    if err != nil {
        return err
    }
    // inspect event
}

Принцип работы

  1. Потребитель вызывает WatchStream, который управляет циклом выполнения.
  2. Суперстеп выполняется до полного завершения, и события накапливаются.
  3. Накопленные события передаются потребителю.
  4. Следующий супершаг начинается только после того, как потребитель получит события предыдущего супершагa.

Когда следует использовать Lockstep

Используйте Lockstep, если детерминированное поведение имеет значение больше, чем потоковая передача с низкой задержкой, например модульные тесты, отладка или сценарии, в которых необходимо полностью обработать события суперстепа перед началом следующего суперстепа.

Выбор режима выполнения

Для большинства рабочих сценариев используйте inproc.Default или inproc.OffThread. Используйте inproc.Lockstep , если детерминированное упорядочивание событий является более важным, чем задержка потоковой передачи, например в тестах. Используйте inproc.Concurrent только в том случае, если каждая привязка в рабочем процессе поддерживает параллельное общее выполнение.

// Production: OffThread (default)
stream, err := inproc.Default.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer stream.Close(ctx)

// Testing: Lockstep for deterministic behavior
testStream, err := inproc.Lockstep.RunStreaming(ctx, wf, input)
if err != nil {
    return err
}
defer testStream.Close(ctx)

Операция без потокового режима

Во всех средах выполнения также поддерживается непотоковый режим Run, который выполняется до следующей остановки, а созданные события сохраняются в возвращаемом объекте run.

run, err := inproc.Default.Run(ctx, wf, input)
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)
    }
}

Дальнейшие действия