For AI agents: the complete documentation index is available at https://a3s-lab.github.io/Flow/llms.txt, the full documentation bundle is available at https://a3s-lab.github.io/Flow/llms-full.txt, and this page is available as Markdown at https://a3s-lab.github.io/Flow/guide/retries-and-waits.md.
  • 简体中文
  • v1.1.0
  • 重试与等待

    步骤失败和时间等待都会让运行暂时停下来。两者的持久状态不同。重试属于某个步骤尝试,等待是工作流显式声明的 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 包含第一次执行。上面的策略最多执行四次,失败后等待五秒再开始下一次。

    三种常用策略

    构造方法行为适合场景
    RetryPolicy::none()只执行一次业务拒绝、确定性校验失败
    RetryPolicy::fixed(n, delay)每次使用同一延迟很短的服务抖动、固定轮询间隔
    RetryPolicy::exponential(n, initial, max)指数上限加确定性 full jitter限流、共享依赖故障、较长恢复窗口

    指数策略从 1..=current_cap 毫秒选择延迟。选择依据是不可变的运行 ID、步骤 ID 和失败尝试号。进程重启后会得到同一个延迟,不需要保存随机数生成器状态。

    let retry = RetryPolicy::exponential(
        8,
        Duration::from_secs(1),
        Duration::from_secs(30),
    );

    初始延迟至少是一毫秒,最大延迟不会小于初始延迟。构造方法会把 max_attempts 收紧到至少一次。

    步骤如果执行了较长时间才失败,Flow 会以实际失败时刻而不是本次驱动 开始时刻计算新的 retry_after。调度器提供的未来截止时间仍会作为下界保留, 因此补偿式追赶不会让长时间运行的尝试立刻过期。

    重试耗尽后继续工作流

    默认行为是 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。 对于仍在等待且尚未到期的定时器,直接调用 resume_wait() 会返回状态迁移错误; 已经终止的运行再次投递仍是幂等空操作。

    调度器怎样唤醒运行

    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

    无论选择哪种方式,外部步骤仍是至少一次交付。重试次数不能替代幂等设计。