重试与等待
步骤失败和时间等待都会让运行暂时停下来。两者的持久状态不同。重试属于某个步骤尝试,等待是工作流显式声明的 UTC 截止时间。它们都能释放当前 worker,并由调度器在到期后继续。
默认重试策略
ctx.schedule_step() 使用 RetryPolicy::default()。默认最多执行三次,尝试之间没有延迟,最后一次仍失败时终止运行。
生产代码通常应该显式写出策略。这样读工作流的人不用去查默认值,历史中的意图也更清楚。
use a3s_flow::RetryPolicy;
use std::time::Duration;
let retry = RetryPolicy::fixed(4, Duration::from_secs(5));
Ok(ctx.schedule_step_with_retry(
"capture-payment",
"capture_payment",
input,
retry,
))
max_attempts 包含第一次执行。上面的策略最多执行四次,失败后等待五秒再开始下一次。
三种常用策略
指数策略从 1..=current_cap 毫秒选择延迟。选择依据是不可变的运行 ID、步骤 ID 和失败尝试号。进程重启后会得到同一个延迟,不需要保存随机数生成器状态。
let retry = RetryPolicy::exponential(
8,
Duration::from_secs(1),
Duration::from_secs(30),
);
初始延迟至少是一毫秒,最大延迟不会小于初始延迟。构造方法会把 max_attempts 收紧到至少一次。
重试耗尽后继续工作流
默认行为是 StepFailureAction::FailRun。有些流程需要在最后一次失败后进入人工处理、降级路径或补偿逻辑,可以使用 continue_workflow_on_failure()。
let retry = RetryPolicy::fixed(3, Duration::from_secs(2))
.continue_workflow_on_failure();
if let Some(error) = ctx.step_failed("reserve-stock") {
return Ok(ctx.schedule_step(
"open-manual-case",
"open_manual_case",
json!({ "reason": error }),
));
}
Ok(ctx.schedule_step_with_retry(
"reserve-stock",
"reserve_stock",
input,
retry,
))
ctx.step_failed() 只读取已经持久化的最终失败。不要在工作流里根据临时错误字符串猜测当前尝试是否结束。
定时等待
wait_until() 把稳定 wait_id 和绝对 UTC 时间写入历史。到期以前,运行状态为 Suspended,不需要保留异步调用栈。
use chrono::{Duration, Utc};
let resume_at = Utc::now() + Duration::hours(24);
if ctx.wait_completed("payment-window") {
return Ok(ctx.fail("payment window expired"));
}
Ok(ctx.wait_until("payment-window", resume_at))
这个片段中的 Utc::now() 只能在首次产生命令的确定性边界外计算并固定。更稳妥的做法是把截止时间放进运行输入,或由创建运行的主机计算后传入。
let deadline = ctx.input()["payment_deadline"]
.as_str()
.ok_or_else(|| FlowError::Runtime("missing payment_deadline".into()))?
.parse::<chrono::DateTime<Utc>>()?;
Ok(ctx.wait_until("payment-window", deadline))
同一个 wait_id 如果在重放时带来不同截止时间,Flow 返回 NonDeterministic。
调度器怎样唤醒运行
FlowScheduler 查询 FlowEventStore::list_due_wakeups(),其中包括到期等待和延迟重试。内存与 JSONL 存储通过投影历史查询,SQLite 与 PostgreSQL 使用迁移生成的索引投影。
let tick = scheduler.enqueue_due_work(chrono::Utc::now()).await?;
println!("enqueued={}", tick.enqueued_tasks);
next_scheduled_wakeup() 返回整个存储中最早的截止时间,适合让主机决定下一次睡眠。调度循环仍要有退出信号和有界等待,不能用无限 sleep 阻止关闭。
选择重试还是工作流循环
短暂技术故障使用步骤重试。每次尝试仍是同一个步骤身份,重试策略负责次数和延迟。
业务轮询使用工作流循环。每轮有新的稳定步骤 ID,工作流可以检查结果、记录进度,并决定是否继续等待。仓库中的 polling_loop 示例展示了这条路径。
cargo run --example retry_backoff
cargo run --example recoverable_step_failure
cargo run --example scheduler_worker
cargo run --example polling_loop
无论选择哪种方式,外部步骤仍是至少一次交付。重试次数不能替代幂等设计。