fix(flow): 修复流程执行异常处理机制

- 在FlowItemExecutor中添加GlobalException导入并修改分支匹配失败处理逻辑
- 将分支节点未匹配到路径的警告改为抛出异常
- 在FlowTaskAsyncDispatcher中添加异常处理相关依赖注入
- 实现更完善的异常捕获和任务失败标记机制
- 优化中断异常处理确保线程状态正确维护
- 添加任务上下文检查避免重复失败标记
This commit is contained in:
lixiaolong 2026-08-14 15:12:11 +08:00
parent c4a1266335
commit 861252f066
2 changed files with 19 additions and 4 deletions

View File

@ -3,6 +3,7 @@ package com.cmvr.test.flow.runtime.engine;
import cn.hutool.core.bean.BeanUtil;
import cn.hutool.core.util.ObjUtil;
import com.alibaba.fastjson2.JSONObject;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.test.enums.NodeTypeEnum;
import com.cmvr.test.flow.builder.FlowEdge;
import com.cmvr.test.flow.builder.FlowGraph;
@ -260,8 +261,7 @@ public class FlowItemExecutor {
FlowEdge matchedEdge = graph.getEdgeByFromAnchorId(matchedAnchorId);
if (matchedEdge == null) {
log.warn("未找到匹配的分支出边,anchorId={}, nodeId={}", matchedAnchorId, nodeId);
return;
throw new GlobalException("分支节点未匹配到可执行路径,节点: {},锚点: {}", node.getNodeName(), matchedAnchorId);
}
// 移除未命中的路径

View File

@ -1,6 +1,10 @@
package com.cmvr.test.flow.runtime.engine;
import cn.hutool.core.exceptions.ExceptionUtil;
import com.cmvr.test.enums.ActionEnum;
import com.cmvr.test.enums.RunModeEnum;
import com.cmvr.test.flow.context.TaskContext;
import com.cmvr.test.flow.context.TaskInstHolder;
import com.cmvr.test.model.dto.TeQueryTaskDetailDTO;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
@ -22,6 +26,7 @@ public class FlowTaskAsyncDispatcher {
private final ThreadPoolTaskExecutor executor;
private final FlowTaskEngine taskEngine;
private final TaskInstHolder taskInstHolder;
/**
* 异步提交流程任务
@ -42,8 +47,18 @@ public class FlowTaskAsyncDispatcher {
log.info("异步任务执行结束:threadName={}, instId={}, taskId={}", Thread.currentThread().getName(), instId, taskId);
} catch (Exception e) {
log.error("任务执行异常:threadName={}, instId={}, taskId={}, 错误={}", Thread.currentThread().getName(), instId, taskId, e.getMessage(), e);
Thread.currentThread().interrupt(); // 保持中断状态
if (ExceptionUtil.getRootCause(e) instanceof InterruptedException) {
Thread.currentThread().interrupt();
return;
}
TaskContext context = taskInstHolder.getContext(instId);
if (context != null && !context.isStopped()) {
taskInstHolder.markFailed(
instId, taskId, context.getItemId(), null, null,
ActionEnum.NONE.getOperate(), ActionEnum.NONE.getAction(), null,
"任务执行异常: " + e.getMessage(), null);
}
}
});
}
}
}