跳到主要内容

Cron schedule

在一个 durable Flow 里跑固定间隔的工作。

Flow definition

Definition graph

CronScheduleFlow

Valid

python · examples/python/dex_examples/patterns/cron/cron_schedule_flow.py

这个 Flow 有三个 Step definition:

  1. InitStep 校验 interval 和 run count 都是正数,然后转到 WaitForSchedule
  2. WaitForSchedule 等待 interval TimerTriggerSkipSkip 消耗一次 occurrence 并开始下一个 interval。Timer 或 Trigger 会同时启动 Run 和下一次 WaitForSchedule
  3. Run 执行 scheduled work。除非这是最后一次 occurrence,否则它到达 dead end;最后一次会完成 Flow。

外部代码通过 Client 发布 TriggerSkip Channel。

核心实现

每个 sample 都展示 WaitForSchedule Step definition 和它的 scheduling decision。完整 runnable example 还包括 input type、Flow registration 和 work implementation。

class _WaitForSchedule(Step[_ScheduleState]):
def __init__(
self,
trigger: Channel[None],
skip: Channel[None],
) -> None:
self.trigger = trigger
self.skip = skip

def wait_for(self, context: Context, state: _ScheduleState) -> Wait:
return Wait.any_of(
Timer.by_duration(state.interval.duration()),
self.trigger.for_one(),
self.skip.for_one(),
)

def execute(self, context: Context, state: _ScheduleState) -> StepDecision:
if self.skip.results(context):
return self._next_schedule(state)
run_input = _RunInput(
run_number=state.remaining_runs,
is_final=state.remaining_runs == 1,
)
if run_input.is_final:
return go_to(_Run, run_input)
return go_to_many(
StepMovement.of(_Run, run_input),
StepMovement.of(
_WaitForSchedule,
_ScheduleState(state.interval, state.remaining_runs - 1),
),
)

例子: examples/python/dex_examples/patterns/cron/cron_schedule_flow.py

相关页面