fix(flow): 修复任务暂停终止时的节点执行异常处理

- 在 FlowAfterInterceptor 中提前检查任务暂停终止状态并返回 null
- 将 TaskContext 获取移至条件判断前避免空指针异常
- 在 FlowExceptionInterceptor 中处理中断异常和暂停终止状态
- 添加对空节点动作和不支持操作类型的校验
- 在 FlowItemExecutor 中增加任务状态检查防止继续执行
- 实现 handleUnsuccessfulResult 方法统一处理执行失败情况
- 添加 notifyFinished 方法确保流程等待器正确释放
- 在 FlowTaskEngine 中增加任务状态检查和资源清理
- 优化 TaskInstHolder 中的状态同步逻辑
- 添加 FlowFunctionNodeHandler 单元测试验证异常处理逻辑
This commit is contained in:
lixiaolong 2026-08-07 16:48:58 +08:00
parent 54df691ea6
commit c51823801d
8 changed files with 216 additions and 62 deletions

View File

@ -118,12 +118,11 @@ public class TaskInstHolder {
String action, String params, String message, List<Integer> iterations) {
nodeInstService.logFailed(instId, taskId, itemId, nodeId, nodeType, operate, action, params, message,iterations);
syncStatus(instId, TaskStatusEnum.FAILED);
// 先获取上下文在注销前
// 终止和并发失败都可能先注销上下文失败状态仍需落库
TaskContext ctx = taskContextManager.get(instId);
String moduleCode = ctx != null ? ctx.getModuleCode() : null;
syncStatus(instId, TaskStatusEnum.FAILED);
// 注销任务上下文
taskContextManager.unregister(instId);
@ -141,22 +140,27 @@ public class TaskInstHolder {
public void syncStatus(String instId, TaskStatusEnum status) {
TaskContext ctx = taskContextManager.get(instId);
switch (status) {
case PAUSED:
ctx.setPaused(true);
ctx.setStatus(TaskStatusEnum.PAUSED);
ctx.setTerminalStatus(TerminalStatusEnum.PAUSED);
break;
case STOPPED:
ctx.setStopped(true);
ctx.setPaused(false);
ctx.setTerminalLocked(false);
break;
if (ctx != null) {
switch (status) {
case PAUSED:
ctx.setPaused(true);
ctx.setStatus(TaskStatusEnum.PAUSED);
ctx.setTerminalStatus(TerminalStatusEnum.PAUSED);
break;
case STOPPED:
ctx.setStopped(true);
ctx.setPaused(false);
ctx.setTerminalLocked(false);
break;
}
ctx.setStatus(status);
ctx.updateTimestamp();
}
ctx.setStatus(status);
ctx.updateTimestamp();
taskInstService.updateStatus(instId, status);
if (ctx == null) {
return;
}
// 发布任务状态变更事件暂停/终止/恢复
switch (status) {
case PAUSED:

View File

@ -1,6 +1,7 @@
package com.cmvr.test.flow.runtime.dispatcher;
import com.cmvr.test.enums.ActionEnum;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage;
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult;
import com.cmvr.test.flow.runtime.operator.NodeOperateHandler;
@ -20,10 +21,16 @@ public class FlowFunctionNodeHandler implements FlowNodeTypeHandler {
@Override
public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) {
ActionEnum action = message.getAction();
if (action == null) {
throw new GlobalException("节点动作不能为空");
}
NodeOperateHandler handler = operateHandlerMap.get(action.getOperate());
if (handler == null) {
throw new GlobalException("不支持的节点操作类型: " + action.getOperate());
}
TaskNodeExecuteResult execute = handler.execute(message);
return TaskNodeExecuteResult.success(execute.getOutputParams());
// 暂停或终止时操作处理器允许返回 null由执行器统一结束本次调度
return handler.execute(message);
}
}

View File

@ -58,7 +58,12 @@ public class FlowItemExecutor {
executeNode(graph, graph.getStartNode(), startNodeId, startInputDefs, rootMessage, onFinished, new ArrayList<>());
// 调度后续节点并发
flowTaskScheduler.start(graph, taskInstHolder.getContext(rootMessage.getInstId()), executor, onFinished);
TaskContext context = taskInstHolder.getContext(rootMessage.getInstId());
if (context == null || context.isStopped()) {
notifyFinished(onFinished, rootMessage.getInstId(), startNodeId);
return;
}
flowTaskScheduler.start(graph, context, executor, onFinished);
}
/**
@ -86,8 +91,12 @@ public class FlowItemExecutor {
rootMessage, onFinished, iterations);
log.info("当前执行-子图");
// 启动子图调度
flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()),
executor, onFinished, iterations);
TaskContext context = taskInstHolder.getContext(rootMessage.getInstId());
if (context == null || context.isStopped()) {
notifyFinished(onFinished, rootMessage.getInstId(), subStartNodeId);
return;
}
flowTaskScheduler.subStart(subGraph, context, executor, onFinished, iterations);
}
/**
@ -120,7 +129,7 @@ public class FlowItemExecutor {
TaskNodeExecuteResult result = chainBuilder.build(interceptors, dispatchNodeHandler).apply(message);
if (ObjUtil.isEmpty(result) || !result.isSuccess()) {
// log.warn("节点执行失败中断调度nodeId={}, status={}", nodeId, result != null ? result.isSuccess() : "null");
handleUnsuccessfulResult(message, result, onFinished);
return;
}
@ -147,6 +156,20 @@ public class FlowItemExecutor {
n -> executeNode(graph, n, startNodeId, startInputDefs, rootMessage, onFinished, iterations),
iterations
);
} catch (RuntimeException e) {
TaskContext ctx = taskInstHolder.getContext(instId);
try {
if (ctx != null && !ctx.isPaused() && !ctx.isStopped()) {
taskInstHolder.markFailed(
instId, rootMessage.getTaskId(), itemId, nodeId, node.getNodeType().name(),
node.getAction() == null ? null : node.getAction().getOperate(),
node.getAction() == null ? null : node.getAction().getAction(),
null, "节点执行异常: " + e.getMessage(), iterations);
}
} finally {
notifyFinished(onFinished, instId, nodeId);
}
throw e;
} finally {
// 如果不是暂停全中断
TaskContext ctx = taskInstHolder.getContext(instId);
@ -156,6 +179,52 @@ public class FlowItemExecutor {
}
}
private void handleUnsuccessfulResult(TaskNodeExecuteMessage message, TaskNodeExecuteResult result,
Runnable onFinished) {
TaskContext ctx = taskInstHolder.getContext(message.getInstId());
if (result == null) {
// 暂停需要保留等待关系以便恢复终止会注销上下文必须释放等待线程
if (ctx == null || ctx.isStopped()) {
notifyFinished(onFinished, message.getInstId(), message.getNodeId());
} else if (!ctx.isPaused()) {
try {
taskInstHolder.markFailed(
message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(),
message.getNodeType(),
message.getAction() == null ? null : message.getAction().getOperate(),
message.getAction() == null ? null : message.getAction().getAction(),
message.getInputParams() == null ? null : message.getInputParams().toJSONString(),
"节点执行结果为空", message.getIterations());
} finally {
notifyFinished(onFinished, message.getInstId(), message.getNodeId());
}
}
return;
}
try {
if (ctx != null && !ctx.isPaused() && !ctx.isStopped()) {
taskInstHolder.markFailed(
message.getInstId(), message.getTaskId(), message.getItemId(), message.getNodeId(),
message.getNodeType(),
message.getAction() == null ? null : message.getAction().getOperate(),
message.getAction() == null ? null : message.getAction().getAction(),
message.getInputParams() == null ? null : message.getInputParams().toJSONString(),
result.getErrorMsg(), message.getIterations());
}
} finally {
notifyFinished(onFinished, message.getInstId(), message.getNodeId());
}
}
private void notifyFinished(Runnable onFinished, String instId, String nodeId) {
try {
onFinished.run();
} catch (Exception e) {
log.error("释放流程等待器失败: instId={}, nodeId={}", instId, nodeId, e);
}
}
@NotNull
private static TaskNodeExecuteMessage buildTaskNodeExecuteMsg(FlowGraph graph, FlowNodeWrapper node,
TaskNodeExecuteMessage rootMessage, String nodeId, String nodeName,

View File

@ -22,45 +22,60 @@ public class FlowTaskEngine {
private final FlowItemExecutor flowTaskExecutor;
private final FlowTaskScheduler flowTaskScheduler;
public void run(String instId, String taskId, String robotId,
RunModeEnum runMode, List<TeQueryTaskDetailDTO> items) {
TaskContext context = taskInstHolder.getContext(instId);
context.setRunMode(runMode);
JSONObject runParams = context.getRunParams();
int totalItemCount = items.size();
for (int i = 0; i < totalItemCount; i++) {
TeQueryTaskDetailDTO item = items.get(i);
// 等待当前 item 完成后再执行下一个
CountDownLatch latch = new CountDownLatch(1);
String itemId = item.getDetectItemId();
context.setItemId(itemId);
TaskNodeExecuteMessage message = new TaskNodeExecuteMessage();
message.setInstId(instId);
message.setTaskId(taskId);
message.setRobotId(robotId);
message.setItemId(itemId);
if (runMode.equals(RunModeEnum.VI_PROJECT)) {
message.setLoopDetail(item.getSchemeInfo());
try {
TaskContext context = taskInstHolder.getContext(instId);
if (context == null) {
log.info("任务已终止,跳过执行: instId={}", instId);
return;
}
context.setRunMode(runMode);
JSONObject runParams = context.getRunParams();
int pendingItemCount = totalItemCount - i - 1;
message.setPendingItemCount(pendingItemCount);
int totalItemCount = items.size();
JSONObject jsonObject = runParams.getJSONObject(itemId);
if (!runMode.equals(RunModeEnum.TRIAL)) {
context.setRunParams(jsonObject);
}
flowTaskExecutor.executeGraph(message, item, latch::countDown);
try {
latch.await(); // 等待该检测项的子流程全部完成
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.error("检测项执行被中断itemId={}", itemId, e);
break;
for (int i = 0; i < totalItemCount; i++) {
TeQueryTaskDetailDTO item = items.get(i);
// 等待当前 item 完成后再执行下一个
CountDownLatch latch = new CountDownLatch(1);
String itemId = item.getDetectItemId();
context.setItemId(itemId);
TaskNodeExecuteMessage message = new TaskNodeExecuteMessage();
message.setInstId(instId);
message.setTaskId(taskId);
message.setRobotId(robotId);
message.setItemId(itemId);
if (runMode.equals(RunModeEnum.VI_PROJECT)) {
message.setLoopDetail(item.getSchemeInfo());
}
int pendingItemCount = totalItemCount - i - 1;
message.setPendingItemCount(pendingItemCount);
JSONObject jsonObject = runParams == null ? null : runParams.getJSONObject(itemId);
if (!runMode.equals(RunModeEnum.TRIAL)) {
context.setRunParams(jsonObject);
}
flowTaskExecutor.executeGraph(message, item, latch::countDown);
try {
latch.await(); // 等待该检测项的子流程全部完成
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
log.info("检测项执行被终止itemId={}", itemId);
break;
}
context = taskInstHolder.getContext(instId);
if (context == null || context.isStopped()) {
break;
}
}
} finally {
flowTaskScheduler.clear(instId);
}
}
}

View File

@ -28,6 +28,11 @@ public class FlowAfterInterceptor extends AbstractFlowMsgPreInterceptor{
@Override
public TaskNodeExecuteResult doIntercept(TaskNodeExecuteMessage message, Function<TaskNodeExecuteMessage, TaskNodeExecuteResult> next) {
TaskNodeExecuteResult result = next.apply(message);
TaskContext taskContext = taskInstHolder.getContext(message.getInstId());
if (result == null && (taskContext == null || taskContext.isPaused() || taskContext.isStopped())) {
return null;
}
FlowExecutionEvent.EventType eventType;
String errorMessage = null;
@ -38,7 +43,6 @@ public class FlowAfterInterceptor extends AbstractFlowMsgPreInterceptor{
errorMessage = result != null ? result.getErrorMsg() : "执行结果为空";
}
TaskContext taskContext = taskInstHolder.getContext(message.getInstId());
String moduleCode = taskContext != null ? taskContext.getModuleCode() : null;
if (moduleCode == null) {

View File

@ -4,6 +4,7 @@ import cn.hutool.core.exceptions.ExceptionUtil;
import cn.hutool.core.util.ObjUtil;
import cn.hutool.core.util.StrUtil;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.test.flow.context.TaskContext;
import com.cmvr.test.flow.context.TaskInstHolder;
import com.cmvr.test.flow.runtime.event.FlowExecutionEventPublisher;
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage;
@ -32,13 +33,20 @@ public class FlowExceptionInterceptor extends AbstractFlowMsgPreInterceptor {
try {
return next.apply(message);
} catch (Exception e) {
log.error("节点执行异常instId={}, nodeId={}, action={}",
message.getInstId(), message.getNodeId(), message.getAction(), e);
if (ExceptionUtil.getRootCause(e) instanceof InterruptedException) {
log.info("节点被中断退出: {}", message.getNodeId());
log.info("节点因任务控制被中断: instId={}, nodeId={}", message.getInstId(), message.getNodeId());
return null;
}
TaskContext context = taskInstHolder.getContext(message.getInstId());
if (context == null || context.isPaused() || context.isStopped()) {
log.info("任务已暂停或终止,忽略节点取消异常: instId={}, nodeId={}, error={}",
message.getInstId(), message.getNodeId(), ExceptionUtil.getMessage(e));
return null;
}
log.error("节点执行异常instId={}, nodeId={}, action={}",
message.getInstId(), message.getNodeId(), message.getAction(), e);
taskInstHolder.markFailed(
message.getInstId(),
message.getTaskId(),

View File

@ -48,8 +48,13 @@ public class FlowLoggingInterceptor extends AbstractFlowMsgPreInterceptor {
return null;
}
if (result == null || !result.isSuccess()) {
return result;
}
taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(),
result.getOutputParams().toJSONString(), StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations());
result.getOutputParams() == null ? null : result.getOutputParams().toJSONString(),
StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations());
long end = System.currentTimeMillis();
log.info("[ {} ]节点执行结束, 参数[ {} ], 循环次数[ {} ], 耗时[ {} ]", message.getAction(), message.getInputParams(), message.getLoopNum(), end - start);

View File

@ -0,0 +1,42 @@
package com.cmvr.test.flow.runtime.dispatcher;
import com.cmvr.test.enums.ActionEnum;
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteMessage;
import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult;
import com.cmvr.test.flow.runtime.operator.NodeOperateHandler;
import org.junit.Test;
import java.util.Collections;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNull;
public class FlowFunctionNodeHandlerTest {
@Test
public void stoppedNodeMayReturnNullWithoutCausingAnotherException() {
NodeOperateHandler stoppedHandler = message -> null;
FlowFunctionNodeHandler handler = new FlowFunctionNodeHandler(
Collections.singletonMap("EDGE", stoppedHandler));
TaskNodeExecuteMessage message = new TaskNodeExecuteMessage();
message.setAction(ActionEnum.CAMERA_START);
assertNull(handler.handle(message));
}
@Test
public void operationFailureIsNotConvertedToSuccess() {
NodeOperateHandler failedHandler = message -> TaskNodeExecuteResult.failure("device failed");
FlowFunctionNodeHandler handler = new FlowFunctionNodeHandler(
Collections.singletonMap("EDGE", failedHandler));
TaskNodeExecuteMessage message = new TaskNodeExecuteMessage();
message.setAction(ActionEnum.CAMERA_START);
TaskNodeExecuteResult result = handler.handle(message);
assertFalse(result.isSuccess());
assertEquals("device failed", result.getErrorMsg());
}
}