工作流执行模式

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

Overview

OffThread (默认值) Lockstep
超步执行 后台线程 使用者的线程
事件传递 即时,当事件被引发时 在每个超级步骤完成之后进行批处理。
步骤执行 独立于事件处理 暂停,直到批处理事件被处理完毕
并发 使用者在超级步骤运行时读取事件 用户与超级步骤执行交替
最适用于 实时流媒体、生产场景 测试、调试、确定性排序

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

Important

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

Lockstep

在 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

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

OffThread

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

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

工作原理

  1. 后台 goroutine 会在有待处理消息时运行 superstep。
  2. 当执行程序生成输出或事件时,工作流事件将写入流。
  3. 消费者通过 WatchStream 读取事件,并在事件产生时接收它们。
  4. 当所有超级步骤均已完成且没有剩余消息时,此次运行将停止,状态为“空闲”或“待处理请求”。

并发运行

当工作流中的所有执行程序绑定都支持并发共享执行时使用 inproc.Concurrent

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

Lockstep

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

后续步骤