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 的异常处理机制。
说到底,工作流引擎真正难的地方,从来不是让节点“跑起来”。
而是:
当节点没有按照预期运行时,系统依然知道下一步该怎么办。