From 3d9049bc2492b65a7615216bb1f367b1b89172c9 Mon Sep 17 00:00:00 2001 From: gitjxm <41097032+gitjxm@users.noreply.github.com> Date: Fri, 17 May 2024 10:42:43 +0800 Subject: [PATCH] Fix memory overflow issue caused by workflow circular references When the workflow references in a loop, going online can cause memory overflow. avoid this situation, add access records --- .../service/process/ProcessServiceImpl.java | 19 ++++++++++++++++++- 1 file changed, 18 insertions(+), 1 deletion(-) diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java index 635948fc86..507788d1c3 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessServiceImpl.java @@ -492,10 +492,22 @@ public class ProcessServiceImpl implements ProcessService { */ @Override public List findAllSubWorkflowDefinitionCode(long parentCode) { + return findAllSubWorkflowDefinitionCode(parentCode, new HashSet<>()); + } + + /** + * recursive query sub process definition id by parent id. + * + * @param parentCode parentCode + * @param visited visited + */ + private List findAllSubWorkflowDefinitionCode(long parentCode, Set visited) { List taskNodeList = taskDefinitionDao.getTaskDefinitionListByDefinition(parentCode); if (CollectionUtils.isEmpty(taskNodeList)) { return Collections.emptyList(); } + visited.add(parentCode); + List subWorkflowDefinitionCodes = new ArrayList<>(); for (TaskDefinition taskNode : taskNodeList) { @@ -504,10 +516,15 @@ public class ProcessServiceImpl implements ProcessService { if (parameterJson.get(CMD_PARAM_SUB_PROCESS_DEFINE_CODE) != null) { SubProcessParameters subProcessParam = JSONUtils.parseObject(parameter, SubProcessParameters.class); long subWorkflowDefinitionCode = subProcessParam.getProcessDefinitionCode(); + if (visited.contains(subWorkflowDefinitionCode)) { + throw new ServiceException("Detected a cycle in the workflow definitions with code: " + subWorkflowDefinitionCode); + } subWorkflowDefinitionCodes.add(subWorkflowDefinitionCode); - subWorkflowDefinitionCodes.addAll(findAllSubWorkflowDefinitionCode(subWorkflowDefinitionCode)); + subWorkflowDefinitionCodes.addAll(findAllSubWorkflowDefinitionCode(subWorkflowDefinitionCode, visited)); } } + + visited.remove(parentCode); return subWorkflowDefinitionCodes; }