工作流执行模式

在 .NET 中运行工作流时, 执行模式 控制如何处理超级步骤以及如何将事件传递到使用者。 该 InProcessExecution 类公开两种执行模式: OffThreadLockstep

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. 当所有超级步骤完成且没有消息保留时,运行将停止,并显示IdlePendingRequests状态。

因为超级步循环和使用者并发运行,事件在引发时立即出现 — 没有缓冲延迟。 这使 OffThread 非常适合低延迟事件传递很重要的流式场景,例如在 UI 中逐令牌显示更新。

并发运行

OffThread 还 支持并发变体 ,允许多个运行同时共享同一工作流实例:

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

重要

并发执行要求将工作流中的所有执行程序声明 crossRunShareable (在构造函数上)或作为工厂方法提供。

锁步

在 Lockstep 模式下,超步在 使用者的线程 而不是后台任务上运行。 事件在每个超级步骤期间累积,并在超级步骤完成后作为批次发出。

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

工作原理

  1. 使用者调用 WatchStreamAsync,这将驱动执行循环。
  2. 超级步运行至完成,事件在队列中累积。
  3. 超级步完成后,所有排队的事件都交付给使用者。
  4. 仅在使用者接收到上一个超级步的所有事件后,下一个超级步才开始。

这种交替模式意味着使用者和工作流引擎永远不会同时运行。 事件传递过程是确定性的——保证来自一个超级步骤的所有事件在下一个超级步骤的任何事件发生之前全部到达。

何时使用 Lockstep

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 工作流使用单个执行模型,该模型通过异步生成器处理超级步处理和事件传递。 此模型类似于 .NET Lockstep 模式,除非使用者正在主动从生成器获取事件,否则步骤不会前进。

有关运行 Python 工作流的信息,请参阅 工作流生成器和执行

在 Go 中运行工作流时,执行环境控制如何处理超级步骤以及如何将事件传递到使用者。 该workflow/inproc包公开三个环境:Default/OffThreadLockstepConcurrent

Overview

离线程 / 默认 锁步 并发的
超步执行 后台 goroutine 由事件使用者驱动 后台 goroutine
事件传递 立即,在事件引发时 在流被消费时批处理 立即,在事件引发时
最适用于 实时流媒体、生产场景 测试、调试、确定性排序 具有并发安全绑定的共享工作流实例

OffThread

OffThread 是默认执行模式。 这些是等效的:

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

工作原理

  1. 后台 goroutine 在消息待处理时运行超级步。
  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. 下一个超级步骤仅在使用者收到上一个超级步骤的事件之后开始。

何时使用 Lockstep

当确定性行为比低延迟流式处理更重要时,请使用 Lockstep,例如在单元测试、调试,或者希望在下一个超级步骤开始之前先完整处理完当前超级步骤中的所有事件的场景中。

选择执行模式

对于大多数生产方案,请使用 inproc.Defaultinproc.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, 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)
    }
}

后续步骤