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 失败。 |
GracefulComplete 和 ForceComplete 可以返回 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。希望立即做出决定时,使用 ForceComplete 或 ForceFail。
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