4592 字
约 15 分钟
1
PaiFlow 工作流 Agent 是如何做超时控制与异常处理的?

PaiFlow 工作流 Agent 是如何做超时控制与异常处理的?

工作流真正跑到线上以后,最难处理的往往不是“正常执行”,而是各种异常情况。

比如:

  • 大模型接口迟迟不返回;
  • 插件依赖的第三方服务挂了;
  • 网络发生短暂抖动;
  • 节点参数配置错误;
  • 某个节点执行失败;
  • 工作流执行过程中出现程序异常。

如果没有一套完整的异常处理机制,一个节点卡住,就可能拖住整个工作流。

所以在 PaiFlow 里,我们对异常处理做了三个核心设计:

超时控制、失败重试、异常策略。

整个逻辑可以简单理解成:

执行节点
   ↓
是否超时?
   ↓
执行成功? ─────→ 是 → 正常返回
   ↓ 否
是否允许重试?
   ↓
还有重试次数?
   ↓
重新执行
   ↓
仍然失败
   ↓
根据错误策略处理
   ├─ 中断工作流
   ├─ 返回兜底结果继续执行
   └─ 跳转异常分支

这就是 PaiFlow 节点级容错机制的核心。


1. 为什么工作流一定要做超时控制?

超时控制本质上是一种止损机制。

假设工作流中有这样一个节点:

Start
  ↓
调用天气插件
  ↓
LLM 生成旅行方案
  ↓
End

天气插件依赖一个第三方 HTTP 服务。

正常情况下,接口可能 1 秒就返回了。

但如果第三方服务发生故障,HTTP 请求一直阻塞:

调用天气接口
      ↓
一直没有响应
      ↓
工作流一直等待
      ↓
用户看到“执行中……”

如果没有超时机制,这个节点理论上可能一直占用线程。

当这种慢任务越来越多时,还可能进一步造成:

慢任务堆积
   ↓
线程池耗尽
   ↓
新的任务无法执行
   ↓
整个工作流服务被拖慢

因此,一个节点不能无限等待。

它必须有一个明确的时间边界:

超过指定时间仍然没有完成,就认为本次执行失败。


2. RetryConfig:节点容错配置

PaiFlow 把节点的超时和重试配置统一放在了:

RetryConfig

它主要包含几个字段:

shouldRetry
maxRetries
timeout
errorStrategy
customOutput

例如:

{
  "retryConfig": {
    "shouldRetry": true,
    "maxRetries": 1,
    "timeout": 5,
    "errorStrategy": 1,
    "customOutput": {
      "output": "服务暂时不可用"
    }
  }
}

这段配置翻译成人话就是:

允许重试

最多重试 1 次

单次执行超过 5 秒认为超时

最终失败以后采用错误码兜底策略

兜底结果:
"服务暂时不可用"

这样,不同节点就可以拥有不同的容错能力。

例如外部 API 节点可以:

timeout = 10s
retry = 2

而一些不适合重复执行的节点则可以:

retry = false

这也是为什么超时和重试最好做成节点级配置,而不是整个工作流一刀切。


3. AsyncUtil:给节点执行加一个时间边界

PaiFlow 的超时能力主要通过:

AsyncUtil

来实现。

它底层使用的是 Guava 的:

SimpleTimeLimiter

核心代码可以简化成:

public static <T> T callWithTimeLimit(
        long time,
        TimeUnit unit,
        Callable<T> call)
        throws Exception {

    return simpleTimeLimiter.callWithTimeout(
            call,
            time,
            unit
    );
}

使用的时候,把真正的节点执行逻辑封装成一个 Callable

AsyncUtil.callWithTimeLimit(
    timeout,
    TimeUnit.MILLISECONDS,
    () -> doExecute(nodeState)
);

底层其实是一个非常经典的:

ExecutorService
      +
Future

模型。

执行过程大概是:

Callable
   ↓
提交线程池
   ↓
返回 Future
   ↓
Future.get(timeout)

如果任务正常完成:

future.get()
   ↓
获取结果
   ↓
正常返回

如果超过时间还没有完成:

future.get(timeout)
   ↓
TimeoutException

工作流就不会再继续无限等待这个节点。


4. 超时以后,后台任务真的停止了吗?

这里有一个很容易产生误解的地方。

发生超时时,并不意味着:

后台线程一定已经被强制杀死。

SimpleTimeLimiter 在发生超时时,会尝试:

future.cancel(true);

本质上是向执行线程发送一个:

interrupt

中断信号。

如果任务本身支持中断,比如:

Thread.sleep()

或者程序主动检查:

Thread.currentThread().isInterrupted()

那么任务通常可以及时结束。

但如果节点正在进行某些不可中断的操作,比如某些阻塞 IO:

第三方调用
    ↓
卡在底层网络操作

那么即使工作流已经收到 TimeoutException,底层任务仍然有可能继续运行一段时间。

所以这里要区分两个概念:

超时保证的是调用方不再无限等待。

它并不等价于:

百分之百强制终止底层任务。

因此,真正成熟的系统通常还需要在 HTTP 客户端、数据库客户端、模型 SDK 等组件上配置各自的连接超时和读取超时。


5. AbstractNodeExecutor:统一处理超时和重试

PaiFlow 中所有节点执行器,例如:

LLMNodeExecutor

PluginNodeExecutor

CodeNodeExecutor

都不需要自己实现一遍超时和重试逻辑。

公共能力统一放在:

AbstractNodeExecutor

中。

整体结构大概是:

AbstractNodeExecutor.execute()
          ↓
解析 RetryConfig
          ↓
是否允许重试?
          ↓
doExecuteWithTimeout()
          ↓
真正执行节点
          ↓
成功 / 失败

核心逻辑类似:

public NodeRunResult execute(NodeState nodeState) {

    RetryConfig retryConfig =
        nodeState.node()
                 .getData()
                 .getRetryConfig();

    if (retryConfig == null) {
        return doExecute(nodeState);
    }

    if (!retryConfig.getShouldRetry()) {
        return doExecuteWithTimeout(
            nodeState,
            retryConfig
        );
    }

    while (true) {

        NodeRunResult result =
            doExecuteWithTimeout(
                nodeState,
                retryConfig
            );

        if (result.getStatus().isSuccess()) {
            return result;
        }

        // 判断是否超过最大重试次数
        ...
    }
}

这样设计的好处非常明显。

具体节点只需要关注:

这个节点到底要干什么。

例如:

protected NodeRunResult executeNode(...) {

    // 调用大模型

}

至于:

超时

重试

异常处理

状态上报

全部交给父类统一处理。

这就是典型的模板方法设计。


6. doExecuteWithTimeout 是怎么工作的?

真正把节点执行和超时绑定起来的是:

doExecuteWithTimeout

核心逻辑可以简化为:

protected NodeRunResult doExecuteWithTimeout(
        NodeState nodeState,
        RetryConfig retryConfig) {

    if (!retryConfig.timeOutEnabled()) {
        return doExecute(nodeState);
    }

    try {

        return AsyncUtil.callWithTimeLimit(
            retryConfig.toMillis(),
            TimeUnit.MILLISECONDS,
            () -> doExecute(nodeState)
        );

    } catch (TimeoutException e) {

        NodeRunResult result =
            new NodeRunResult();

        result.setError(
            new NodeCustomException(
                ErrorCode.TIMEOUT_ERROR
            )
        );

        return errorResponse(
            nodeState,
            result
        );
    }
}

所以一个节点的实际执行过程就变成:

AbstractNodeExecutor
        ↓
doExecuteWithTimeout
        ↓
AsyncUtil
        ↓
doExecute
        ↓
具体 NodeExecutor

任何节点都自动拥有了超时能力。


7. 为什么失败以后还要重试?

并不是所有失败都意味着真正失败。

很多异常只是暂时的。

比如:

网络短暂波动

DNS 查询失败

第三方接口偶发 502

模型服务瞬时繁忙

连接池暂时耗尽

这种情况下,如果一次失败就直接终止流程,其实有点可惜。

所以 PaiFlow 支持节点失败后进行有限次数重试。

逻辑大概是:

第一次执行
    ↓
失败
    ↓
还有重试次数
    ↓
第二次执行
    ↓
成功
    ↓
继续工作流

对应代码类似:

while (true) {

    NodeRunResult result =
        doExecuteWithTimeout(
            nodeState,
            retryConfig
        );

    if (result.getStatus().isSuccess()) {
        return result;
    }

    if (executeTime >
        retryConfig.getMaxRetries()) {

        return result;
    }

    executeTime =
        node.getExecutedCount()
            .addAndGet(1);
}

通过:

AtomicInteger

记录当前节点已经执行了多少次。

超过:

maxRetries

以后,就停止重试。


8. 重试不是越多越好

重试其实是一把双刃剑。

比如下游服务已经宕机:

第一次失败
 ↓
立即重试
 ↓
再次失败
 ↓
继续重试
 ↓
大量请求一起重试

最终很容易演变成:

重试风暴

下游原本只是短暂故障,结果被大量重试请求彻底压垮。

因此真正成熟的重试策略一般还需要考虑:

最大重试次数

重试间隔

指数退避

随机抖动

异常类型

例如:

第一次失败
等待 1 秒

第二次失败
等待 2 秒

第三次失败
等待 4 秒

而不是:

失败 → 立即重试 → 失败 → 立即重试

PaiFlow 当前实现主要聚焦:

最大重试次数
+
立即重试

但后面完全可以进一步扩展:

retryInterval

exponentialBackoff

jitter

9. 什么错误不应该重试?

这里还有一个非常重要的问题:

是不是所有异常都值得重试?

当然不是。

例如:

参数缺失

JSON 格式错误

节点配置错误

鉴权失败

DSL 不合法

这些都属于:

确定性失败

再执行一次结果基本不会发生变化。

而像:

网络超时

连接失败

服务临时不可用

HTTP 502 / 503

这类错误才比较适合重试。

所以更完整的重试模型应该是:

发生异常
   ↓
判断异常类型
   ↓
是否属于可重试异常?
   ├─ 否 → 直接进入错误处理
   └─ 是 → 判断重试次数

这也是后续可以继续优化的地方。


10. 重试耗尽以后怎么办?

这是异常处理最核心的一层。

假设:

模型调用失败

自动重试 2 次

仍然失败

这时候系统必须回答一个问题:

整个工作流还要不要继续?

不同业务的答案完全不同。

所以 PaiFlow 没有直接:

throw exception;

而是引入了:

ErrorStrategyEnum

目前主要支持三种策略:

INTERUPT

ERR_CODE

ERR_CONDITION

分别代表:

中断

兜底继续

异常分支

11. 策略一:直接中断工作流

第一种:

INTERUPT

是最严格的处理策略。

例如:

Start
 ↓
用户认证
 ↓
支付
 ↓
生成订单
 ↓
End

如果支付节点失败了:

支付失败
   ↓
整个工作流终止

这类关键业务通常不能继续往下执行。

否则可能出现:

支付失败

但是订单却创建成功

所以这种节点通常采用:

失败即中断

12. 策略二:返回兜底结果继续执行

第二种:

ERR_CODE

更像一种:

降级策略

例如:

用户输入
   ↓
查询天气
   ↓
LLM 生成旅行建议
   ↓
End

假设天气接口挂了。

但对于旅行建议来说,天气信息并不是绝对必须。

那么就可以配置:

{
  "customOutput": {
    "weather": "天气数据暂不可用"
  }
}

实际执行就会变成:

天气接口失败
      ↓
重试失败
      ↓
写入 customOutput
      ↓
继续执行 LLM

最终用户可能看到:

天气服务当前不可用,
以下先根据目的地信息提供旅行建议。

整个工作流依然可以完成。

这就是典型的:

服务降级。


13. 策略三:进入异常分支

第三种:

ERR_CONDITION

更像传统工作流引擎中的:

异常边

例如:

调用主模型
    ↓
成功
    ↓
生成结果

如果主模型失败:

调用主模型
    ↓
失败
    ↓
异常分支
    ↓
调用备用模型

也可以设计成:

第三方服务失败
       ↓
异常分支
       ↓
记录错误日志
       ↓
发送告警
       ↓
执行补偿逻辑

到了这个层面以后,异常已经不再只是:

Exception

它本身也成为了:

工作流的一种执行路径。

这也是工作流系统相比普通业务代码更有意思的地方。


14. errorResponse:决定工作流下一步怎么走

PaiFlow 中最终负责这个决策的方法是:

errorResponse

整体逻辑可以简化成:

节点失败
   ↓
读取 RetryConfig
   ↓
读取 ErrorStrategy
   ↓
判断处理策略

如果是:

ERR_CODE

则:

写入 customOutput
       ↓
VariablePool
       ↓
后续节点继续执行

如果是:

ERR_CONDITION

则:

设置 ERR_FAIL_CONDITION
       ↓
工作流调度到异常分支

如果是:

INTERUPT

则:

设置 ERR_INTERUPT
       ↓
发送 onNodeInterrupt
       ↓
结束工作流

所以 errorResponse 本质上不是一个简单的异常转换方法。

它实际上是在做:

异常发生之后的流程决策。


15. NodeCustomException:统一异常模型

节点底层可能抛出各种异常:

TimeoutException

IOException

RuntimeException

第三方 SDK Exception

如果这些异常全部直接暴露给工作流上层,会很难统一处理。

所以 PaiFlow 又增加了一层:

NodeCustomException

例如:

public class NodeCustomException
        extends RuntimeException {

    private final int code;

    private final String msg;
}

通过它统一表达:

错误码

错误信息

错误类型

比如:

TIMEOUT_ERROR

NODE_RUN_ERROR

PARAM_ERROR

上层就不需要了解底层到底发生的是哪种 Java Exception。

它只需要关心:

这个节点发生了什么业务错误?

16. 异常还要告诉前端

异常处理不能只存在于后台日志中。

因为对于工作流系统来说,用户也需要知道:

哪个节点失败了?

为什么失败?

有没有重试?

最终工作流停止了吗?

还是已经走了降级策略?

所以异常发生以后,还会通过:

WorkflowMsgCallback

发送事件。

例如:

onNodeEnd

onNodeInterrupt

再通过上一篇提到的:

WorkflowMsgCallback
       ↓
streamQueue
       ↓
SseStreamCallback
       ↓
SSE
       ↓
前端

实时推送给用户。

所以一个节点超时以后,整个链路其实是:

NodeExecutor
     ↓
执行超过 timeout
     ↓
TimeoutException
     ↓
NodeCustomException
     ↓
errorResponse
     ↓
决定异常策略
     ↓
WorkflowMsgCallback
     ↓
SSE
     ↓
前端展示节点失败状态

这样后台行为和用户看到的状态才能保持一致。


17. 把整个容错链路串起来

到这里,PaiFlow 整个异常处理机制就比较清楚了。

一个节点完整的执行生命周期可以表示成:

                执行节点
                   ↓
             是否配置 timeout
                   ↓
          doExecuteWithTimeout
                   ↓
               执行成功?
              ↙       ↘
            是          否
            ↓           ↓
         正常返回    是否允许重试
                       ↓
                  是否超过次数
                    ↙     ↘
                  否        是
                  ↓         ↓
               再次执行   errorResponse
                            ↓
                     ErrorStrategy
                     ↙    ↓     ↘
               INTERUPT ERR_CODE ERR_CONDITION
                   ↓       ↓        ↓
                 中断    兜底继续   异常分支

这实际上就是 PaiFlow 的:

节点级容错状态机。


18. 一个完整例子

假设现在有一个 LLM 节点:

{
  "retryConfig": {
    "shouldRetry": true,
    "maxRetries": 1,
    "timeout": 5,
    "errorStrategy": 1,
    "customOutput": {
      "output": "AI 服务暂时不可用"
    }
  }
}

第一次执行:

调用 LLM
   ↓
5 秒没有返回
   ↓
TimeoutException

因为配置了:

maxRetries = 1

于是重新执行一次:

第二次调用 LLM
      ↓
依然超时

这时候重试次数耗尽。

进入:

errorResponse

发现配置的是:

ERR_CODE

于是系统将:

{
  "output": "AI 服务暂时不可用"
}

写入:

VariablePool

然后后面的节点依然可以继续执行。

最终工作流不会因为一个非关键节点失败而整体挂掉。


19. 真正需要注意的几个问题

超时、重试这些东西看起来简单,但真正放到生产环境里,需要非常克制。

第一,超时时间不能拍脑袋。

不同节点差异很大:

数据库查询

HTTP 接口

文件处理

大模型调用

正常耗时完全不同。

超时设置得太大:

故障发现太慢

设置得太小:

正常请求也被当成超时

所以最好结合线上:

P95

P99

延迟数据动态调整。

第二,重试一定要考虑幂等性。

例如:

支付

扣款

创建订单

发货

如果接口本身不支持幂等,随便重试可能造成严重的数据问题。

第三,一定要防止重试风暴。

如果下游服务已经挂了:

1000 个请求失败
       ↓
每个请求重试 3 次
       ↓
瞬间变成 4000 次调用

这只会让问题更严重。

所以成熟系统一般还会配合:

指数退避

随机抖动

限流

熔断

一起使用。


20. 小结

PaiFlow 的异常处理,并不是简单地:

try {
    execute();
} catch (Exception e) {
    log.error(...);
}

而是把异常真正纳入了工作流执行模型。

整个体系可以浓缩成四个关键词:

超时
 ↓
重试
 ↓
异常策略
 ↓
状态上报

超时解决:

节点不能无限卡住。

重试解决:

偶发失败可以自动恢复。

错误策略解决:

恢复不了以后,工作流该怎么办。

状态上报解决:

用户和开发者能够知道系统发生了什么。

所以最终 PaiFlow 的节点执行,并不是简单的:

成功 / 失败

而是:

执行
 ↓
超时控制
 ↓
失败判断
 ↓
自动重试
 ↓
异常决策
 ↓
中断 / 降级 / 异常分支
 ↓
状态实时回传

这才是一套真正适合工作流 Agent 的异常处理机制。

说到底,工作流引擎真正难的地方,从来不是让节点“跑起来”。

而是:

当节点没有按照预期运行时,系统依然知道下一步该怎么办。

PaiFlow 工作流 Agent 是如何做超时控制与异常处理的?
http://www.clxhxhhr.top/posts/643/
作者
clxstart
发布于
2026-09-16
许可协议
CC BY-NC-SA 4.0
评论
0 条
还没有评论,先写一条吧。