跳到主要内容

Step decision

本页介绍 StepDecision 的进阶功能。基础功能见基础页面

多个后续 Step

GoTo 启动一个后续 Step。GoToMany 启动多个后续 Step。用它 fan-out 可并行的工作。每个 Step movement 有各自的输入,也可以有各自的 Step options。

return go_to_many(
StepMovement.of(CarrierAStep, Quote(carrier="A", price=10)),
StepMovement.of(CarrierBStep, Quote(carrier="B", price=12)),
StepMovement.of(WinnerStep, quote),
)

例子: examples/python/dex_examples/primitives/step_decision/step_decision_flow.py

关闭决策

决策行为
GracefulComplete停止当前分支,并触发 Flow 为完成做准备。所有分支和 Step 都停止后,Flow 完成。
DeadEnd只停止当前分支。
ForceComplete立即完成 Flow。
ForceFail立即让 Flow 失败。

GracefulCompleteForceComplete 可以返回 output。ForceFail 返回 failure reason。Flow result 会连同 Step type 和 Step execution ID 一起返回每个 output。多个 Step graceful complete 时,它们的所有 output 会一起返回。

Flow 可以没有任何 active Step,但仍继续运行。DeadEnd 停止当前分支,不做出会触发 Flow 完成的决策。

希望剩余 active Step 和分支 graceful finish 时,使用 GracefulComplete。不希望当前分支做出会触发 Flow 完成的决策时,使用 DeadEnd。希望立即做出决定时,使用 ForceCompleteForceFail

if mode == "graceful":
return graceful_complete("done")
if mode == "dead-end":
return go_to_many(
StepMovement.of(BranchWorkerStep, "left"),
StepMovement.of(BranchWorkerStep, "right"),
)

return dead_end()

例子: examples/python/dex_examples/primitives/step_decision/step_decision_flow.py

def execute(self, context: Context, input: None) -> StepDecision:
return force_fail("Workflow did not finish the task in time")

def execute(self, context: Context, input: bool) -> StepDecision:
return force_complete("Workflow completed successfully")

例子: examples/python/dex_examples/patterns/timeout/flow_graceful_timeout.py

Cancellation

返回一个可以取消其他 active Step 的 StepDecision

有两种 cancellation:

  • CancelSteps 按 Step type 选择 active Step。
  • CancelSiblingSteps 选择由同一个 parent Step execution 触发的 active Step。
return (
go_to(RecordQuoteStep, quote)
.with_canceling_steps(self.flow.carrier_a, self.flow.carrier_b)
)

例子: examples/python/dex_examples/primitives/step_decision/step_decision_flow.py

ForceCompleteIfChannelsEmpty

ForceCompleteIfChannelsEmpty 原子检查所有指定的 Channel 是否为空。Channel 中有排队消息时,该 Channel 视为非空。消息从 Channel 消费并提交给 Step 后,在该 Step 完成 Execute 处理之前,该 Channel 也视为非空。若所有指定的 Channel 都为空,它完成 Flow;否则,它转到其他 Step。

当你想 drain Channel 并完成 Flow,让 Flow 保持短小,使用它。完整例子见排空外部 Channel

return force_complete_if_channels_empty(
None,
StepMovement.of(ProcessMessage, None),
self.queue_channel,
)

例子: examples/python/dex_examples/patterns/drain-channels/external_publishing/draining_channel_flow.py