可中断执行
在收到 durable 中断信号后,安全地停止长时间工作。
这个模式适合会长时间循环运行的工作,例如轮询或批处理。工作运行期间,用户可以请求中断。每个分支会完成当前工作单元,在下一轮循环前读取 durable 信号,然后安全地结束。
InterruptibleFlow 会并行启动 WorkAStep 和 WorkBStep。每个 Step 会在安排下一项工作前检查 durable 的 interruptSignal。interrupt RPC 写入该信号,使任一分支都能优雅地完成 Flow。
核心实现
每个示例展示 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