跳到主要内容

可中断执行

在收到 durable 中断信号后,安全地停止长时间工作。

这个模式适合会长时间循环运行的工作,例如轮询或批处理。工作运行期间,用户可以请求中断。每个分支会完成当前工作单元,在下一轮循环前读取 durable 信号,然后安全地结束。

InterruptibleFlow 会并行启动 WorkAStepWorkBStep。每个 Step 会在安排下一项工作前检查 durable 的 interruptSignalinterrupt RPC 写入该信号,使任一分支都能优雅地完成 Flow。

Definition graph

InterruptibleFlow

Valid

python · examples/python/dex_examples/patterns/interruptible/interruptible_execution_flow.py

核心实现

每个示例展示 WorkAStep。它在工作单元之间等待;收到中断信号时完成,否则安排下一轮循环。WorkBStep 采用相同模式。可运行示例还包含 Flow 注册和 interrupt RPC。

class WorkAStep(Step[WorkJobParametersInput]):
def __init__(self, interrupt_signal: Attribute[str]) -> None:
self.interrupt_signal = interrupt_signal

def wait_for(self, context: Context, input: WorkJobParametersInput) -> Wait:
return Wait.until(Timer.by_duration(timedelta(seconds=2)))

def execute(
self,
context: Context,
input: WorkJobParametersInput,
) -> StepDecision:
if (self.interrupt_signal.get(context) or "") == INTERRUPT_VALUE:
print("A: Interrupted!")
return graceful_complete()

if input.progress > input.job_upper_bound:
print("WorkAStep completed")
return graceful_complete()

print(
f"[{context.flow_id}][{context.step_execution_id}]: "
f"Doing job {input.progress}"
)
return go_to(
WorkAStep,
WorkJobParametersInput(input.job_upper_bound, input.progress + 1),
)

例子: examples/python/dex_examples/patterns/interruptible/interruptible_execution_flow.py

相关内容