refactor(flow): 重构流程执行事件和任务状态管理
- 更新 FlowExecutionEvent 类的注释和字段描述,明确流程执行事件的用途和各字段含义 - 重命名 TaskInstHolder.markSuccess 为 markNodeSuccess 方法,并拆分任务完成逻辑 - 添加 markTaskCompleted 私有方法专门处理任务完成状态更新和事件发布 - 保留旧方法名作为兼容性支持,标记为 @Deprecated - 修复 FlowLoggingInterceptor 中的方法调用以使用新的 markNodeSuccess 方法 - 优化 FlowTaskRuntimeEntry 中的终端ID检查逻辑,移除冗余空值判断 - 改进 TaskThreadRegistry.interruptAll 方法,精确匹配实例ID并添加线程中断状态检查 - 在任务中断后释放所有 latch 避免死锁,提升系统稳定性
This commit is contained in:
parent
659073db03
commit
710f33d882
@ -49,30 +49,68 @@ public class TaskInstHolder {
|
||||
nodeInstService.logRunning(instId, taskId, itemId, nodeId, nodeType, operate, action, paramsIn, iterations);
|
||||
}
|
||||
|
||||
/**
|
||||
* 标记节点执行成功(仅记录节点日志,不更新任务状态)
|
||||
*
|
||||
* @param instId 任务实例ID
|
||||
* @param itemId 检测项ID
|
||||
* @param nodeId 节点ID
|
||||
* @param pendingItemCount 剩余待执行检测项数量
|
||||
* @param nodeType 节点类型
|
||||
* @param paramsOut 输出参数
|
||||
* @param message 日志消息
|
||||
* @param iterations 迭代路径
|
||||
*/
|
||||
public void markNodeSuccess(String instId, String itemId, String nodeId, int pendingItemCount, String nodeType,
|
||||
String paramsOut, String message, List<Integer> iterations) {
|
||||
// 仅记录节点执行日志
|
||||
nodeInstService.logSuccess(instId, itemId, nodeId, paramsOut, message, iterations);
|
||||
|
||||
// 如果是最后一个检测项的END节点,标记任务完成
|
||||
if (pendingItemCount == 0 && NodeTypeEnum.END.getCode().equalsIgnoreCase(nodeType)) {
|
||||
markTaskCompleted(instId, itemId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 标记任务执行完成(更新任务状态为SUCCESS,发布完成事件)
|
||||
*
|
||||
* @param instId 任务实例ID
|
||||
* @param itemId 检测项ID
|
||||
*/
|
||||
private void markTaskCompleted(String instId, String itemId) {
|
||||
// 先获取上下文(在注销前)
|
||||
TaskContext ctx = taskContextManager.get(instId);
|
||||
String taskId = ctx != null ? ctx.getTaskId() : null;
|
||||
String moduleCode = ctx != null ? ctx.getModuleCode() : null;
|
||||
|
||||
// 更新任务状态为SUCCESS
|
||||
syncStatus(instId, TaskStatusEnum.SUCCESS);
|
||||
|
||||
// 注销任务上下文(释放终端锁)
|
||||
taskContextManager.unregister(instId);
|
||||
|
||||
// 发布任务完成事件
|
||||
eventPublisher.publishEvent(FlowExecutionEvent.builder()
|
||||
.eventType(FlowExecutionEvent.EventType.TASK_COMPLETED)
|
||||
.instId(instId)
|
||||
.taskId(taskId)
|
||||
.moduleCode(moduleCode)
|
||||
.itemId(itemId)
|
||||
.status(TaskStatusEnum.SUCCESS)
|
||||
.build());
|
||||
|
||||
log.info("任务执行完成: instId={}, taskId={}", instId, taskId);
|
||||
}
|
||||
|
||||
/**
|
||||
* 兼容旧方法名(已废弃,请使用 markNodeSuccess)
|
||||
* @deprecated 使用 markNodeSuccess 替代
|
||||
*/
|
||||
@Deprecated
|
||||
public void markSuccess(String instId, String itemId, String nodeId, int pendingItemCount, String nodeType,
|
||||
String paramsOut, String message, List<Integer> iterations) {
|
||||
nodeInstService.logSuccess(instId, itemId, nodeId, paramsOut, message,iterations);
|
||||
syncStatus(instId, TaskStatusEnum.SUCCESS);
|
||||
// 如果没有下一个检测项要执行 且是end节点 注销任务上下文
|
||||
if (pendingItemCount == 0 && nodeType.equalsIgnoreCase(NodeTypeEnum.END.getCode())) {
|
||||
// 先获取上下文(在注销前)
|
||||
TaskContext ctx = taskContextManager.get(instId);
|
||||
String taskId = ctx != null ? ctx.getTaskId() : null;
|
||||
String moduleCode = ctx != null ? ctx.getModuleCode() : null;
|
||||
|
||||
// 注销任务上下文
|
||||
taskContextManager.unregister(instId);
|
||||
|
||||
// 发布任务完成事件
|
||||
eventPublisher.publishEvent(FlowExecutionEvent.builder()
|
||||
.eventType(FlowExecutionEvent.EventType.TASK_COMPLETED)
|
||||
.instId(instId)
|
||||
.taskId(taskId)
|
||||
.moduleCode(moduleCode)
|
||||
.itemId(itemId)
|
||||
.status(TaskStatusEnum.SUCCESS)
|
||||
.build());
|
||||
}
|
||||
markNodeSuccess(instId, itemId, nodeId, pendingItemCount, nodeType, paramsOut, message, iterations);
|
||||
}
|
||||
|
||||
public void markFailed(String instId, String taskId, String itemId,
|
||||
|
||||
@ -73,18 +73,21 @@ public class TaskThreadRegistry {
|
||||
*/
|
||||
public void interruptAll(String instId) {
|
||||
for (Map.Entry<String, TaskNodeExecuteContext> entry : nodeContextMap.entrySet()) {
|
||||
if (entry.getKey().startsWith(instId)) {
|
||||
if (entry.getKey().startsWith(instId + "_")) {
|
||||
FlowNodeWrapper node = entry.getValue().getNode();
|
||||
if (node.getNodeType().equals(NodeTypeEnum.LOOP)) {
|
||||
log.info("LOOP节点不直接中断");
|
||||
log.info("LOOP节点不直接中断,等待子图完成");
|
||||
continue;
|
||||
}
|
||||
Thread thread = entry.getValue().getThread();
|
||||
log.info("中断线程:{}", thread.getName());
|
||||
thread.interrupt();
|
||||
if (thread != null && !thread.isInterrupted()) {
|
||||
log.info("中断线程:{}", thread.getName());
|
||||
thread.interrupt();
|
||||
}
|
||||
}
|
||||
}
|
||||
// releaseAllLatches(instId);
|
||||
// ✅ 释放所有 latch,避免死锁
|
||||
releaseAllLatches(instId);
|
||||
log.info("-------------------------------------");
|
||||
log.info("pending:\n{}", getPendingNodes(instId).stream().map(x -> x.getNode().getNodeId()).collect(Collectors.toList()));
|
||||
log.info("-------------------------------------");
|
||||
|
||||
@ -61,9 +61,9 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
JSONObject runParams = taskExecuteTrailVO.getRunParams();
|
||||
String terminalId = taskExecuteTrailVO.getTerminalId();
|
||||
if (StrUtil.isEmpty(terminalId)) {
|
||||
// throw new GlobalException("终端ID不能为空");
|
||||
throw new GlobalException("终端ID不能为空");
|
||||
}
|
||||
if (terminalId != null && taskInstHolder.isTerminalLocked(terminalId)) {
|
||||
if (taskInstHolder.isTerminalLocked(terminalId)) {
|
||||
throw new GlobalException("当前终端正在执行其他任务,请稍后再试");
|
||||
}
|
||||
|
||||
@ -124,12 +124,13 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public String executeTaskInternal(String terminalId, String taskId, String moduleCode,
|
||||
RunModeEnum runMode, JSONObject runParams, Object originalVO) {
|
||||
// 终端id可能为空
|
||||
if (StrUtil.isEmpty(terminalId)) {
|
||||
// throw new GlobalException("终端ID不能为空");
|
||||
}
|
||||
// 检查终端状态
|
||||
if (terminalId != null && taskInstHolder.isTerminalLocked(terminalId)) {
|
||||
// throw new GlobalException("当前终端正在执行其他任务,请稍后再试");
|
||||
if (!StrUtil.isEmpty(terminalId) && taskInstHolder.isTerminalLocked(terminalId)) {
|
||||
throw new GlobalException("当前终端正在执行其他任务,请稍后再试");
|
||||
}
|
||||
// 查询任务详情(检测项信息)
|
||||
List<TeQueryTaskDetailDTO> details = taskOrchestrationService.queryTaskDetail(taskId);
|
||||
|
||||
@ -7,9 +7,9 @@ import lombok.Data;
|
||||
import lombok.NoArgsConstructor;
|
||||
|
||||
/**
|
||||
* 流程执行事件
|
||||
* 流程执行事件。
|
||||
* <p>
|
||||
* 携带流程执行过程中的关键信息,供外部模块监听和处理
|
||||
* 用于描述流程引擎执行过程中的关键节点和任务状态变化,供业务模块监听并处理。
|
||||
*/
|
||||
@Data
|
||||
@Builder
|
||||
@ -18,62 +18,62 @@ import lombok.NoArgsConstructor;
|
||||
public class FlowExecutionEvent {
|
||||
|
||||
/**
|
||||
* 事件类型
|
||||
* 事件类型。
|
||||
*/
|
||||
private EventType eventType;
|
||||
|
||||
/**
|
||||
* 任务实例 ID
|
||||
* 流程任务实例ID。
|
||||
*/
|
||||
private String instId;
|
||||
|
||||
/**
|
||||
* 任务 ID
|
||||
* 流程任务ID。
|
||||
*/
|
||||
private String taskId;
|
||||
|
||||
/**
|
||||
* 模块编码
|
||||
* 事件所属业务模块编码,用于监听器路由隔离。
|
||||
*/
|
||||
private String moduleCode;
|
||||
|
||||
/**
|
||||
* 节点数量
|
||||
* 当前检测项内的节点总数。
|
||||
*/
|
||||
private Integer nodeCount;
|
||||
|
||||
/**
|
||||
* 当前检测项 ID
|
||||
* 当前检测项ID。
|
||||
*/
|
||||
private String itemId;
|
||||
|
||||
/**
|
||||
* 节点 ID(节点级事件时有值)
|
||||
* 当前节点ID,节点级事件时有值。
|
||||
*/
|
||||
private String nodeId;
|
||||
|
||||
/**
|
||||
* 节点类型(节点级事件时有值)
|
||||
* 当前节点类型,节点级事件时有值。
|
||||
*/
|
||||
private String nodeType;
|
||||
|
||||
/**
|
||||
* 节点名称
|
||||
* 当前节点名称。
|
||||
*/
|
||||
private String nodeName;
|
||||
|
||||
/**
|
||||
* 任务状态
|
||||
* 当前任务状态。
|
||||
*/
|
||||
private TaskStatusEnum status;
|
||||
|
||||
/**
|
||||
* 错误信息(失败事件时有值)
|
||||
* 错误信息,失败事件时有值。
|
||||
*/
|
||||
private String errorMessage;
|
||||
|
||||
/**
|
||||
* 事件类型枚举
|
||||
* 流程执行事件类型。
|
||||
*/
|
||||
public enum EventType {
|
||||
/** 节点执行完成 */
|
||||
@ -84,9 +84,9 @@ public class FlowExecutionEvent {
|
||||
TASK_COMPLETED,
|
||||
/** 任务执行失败 */
|
||||
TASK_FAILED,
|
||||
/** 任务被暂停 */
|
||||
/** 任务暂停 */
|
||||
TASK_PAUSED,
|
||||
/** 任务被终止 */
|
||||
/** 任务终止 */
|
||||
TASK_STOPPED,
|
||||
/** 任务恢复执行 */
|
||||
TASK_RESUMED,
|
||||
|
||||
@ -48,7 +48,7 @@ public class FlowLoggingInterceptor extends AbstractFlowMsgPreInterceptor {
|
||||
return null;
|
||||
}
|
||||
|
||||
taskInstHolder.markSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(),
|
||||
taskInstHolder.markNodeSuccess(message.getInstId(), message.getItemId(), message.getNodeId(), message.getPendingItemCount(), message.getNodeType(),
|
||||
result.getOutputParams().toJSONString(), StrUtil.format("[{}] 执行成功", message.getAction()), message.getIterations());
|
||||
|
||||
long end = System.currentTimeMillis();
|
||||
|
||||
Loading…
Reference in New Issue
Block a user