This is an automated email from the ASF dual-hosted git repository. vernedeng pushed a commit to branch branch-1.8 in repository https://gitbox.apache.org/repos/asf/inlong.git
commit e96016a04a07f482bb34e1b11fdaaaf30aead637 Author: kipshi <[email protected]> AuthorDate: Thu Jul 13 15:39:05 2023 +0800 [INLONG-8509][Manager] Optimize preProcessTemplateFileTask in AgentServiceImpl (#8510) Co-authored-by: kipshi <[email protected]> --- .../service/core/impl/AgentServiceImpl.java | 53 ++++++++++++---------- .../service/core/impl/AgentServiceTest.java | 6 +-- 2 files changed, 31 insertions(+), 28 deletions(-) diff --git a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java index ad594c3514..174929741e 100644 --- a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java +++ b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/AgentServiceImpl.java @@ -407,30 +407,35 @@ public class AgentServiceImpl implements AgentService { // find those node whose tag match stream_source tag and agent ip match stream_source agent ip List<StreamSourceEntity> sourceEntities = sourceMapper.selectTemplateSourceByCluster(needCopiedStatusList, Lists.newArrayList(SourceType.FILE), agentClusterName); - Set<GroupStatus> noNeedAddTask = Sets.newHashSet(GroupStatus.SUSPENDED, GroupStatus.SUSPENDING); - sourceEntities.forEach(sourceEntity -> { - InlongGroupEntity groupEntity = groupMapper.selectByGroupId(sourceEntity.getInlongGroupId()); - if (groupEntity != null && noNeedAddTask.contains(GroupStatus.forCode(groupEntity.getStatus()))) { - return; - } - StreamSourceEntity subSource = sourceMapper.selectOneByTemplatedIdAndAgentIp(sourceEntity.getId(), - agentIp); - if (subSource == null) { - // if not, clone a subtask for this Agent. - // note: a new source name with random suffix is generated to adhere to the unique constraint - StreamSourceEntity fileEntity = - CommonBeanUtils.copyProperties(sourceEntity, StreamSourceEntity::new); - fileEntity.setSourceName(fileEntity.getSourceName() + "-" - + RandomStringUtils.randomAlphanumeric(10).toLowerCase(Locale.ROOT)); - fileEntity.setTemplateId(sourceEntity.getId()); - fileEntity.setAgentIp(agentIp); - fileEntity.setStatus(SourceStatus.TO_BE_ISSUED_ADD.getCode()); - // create new sub source task - sourceMapper.insert(fileEntity); - LOGGER.info("Transform new template task({}) for agent({}) in cluster({}).", - fileEntity.getId(), taskRequest.getAgentIp(), taskRequest.getClusterName()); - } - }); + Set<GroupStatus> noNeedAddTask = Sets.newHashSet( + GroupStatus.SUSPENDED, GroupStatus.SUSPENDING, GroupStatus.DELETING, GroupStatus.DELETED); + sourceEntities.stream() + .forEach(sourceEntity -> { + InlongGroupEntity groupEntity = groupMapper.selectByGroupId(sourceEntity.getInlongGroupId()); + if (groupEntity != null && noNeedAddTask.contains(GroupStatus.forCode(groupEntity.getStatus()))) { + return; + } + StreamSourceEntity subSource = sourceMapper.selectOneByTemplatedIdAndAgentIp(sourceEntity.getId(), + agentIp); + if (subSource == null) { + InlongClusterNodeEntity clusterNodeEntity = selectByIpAndCluster(agentClusterName, agentIp); + // if stream_source match node_group with node, clone a subtask for this Agent. + // note: a new source name with random suffix is generated to adhere to the unique constraint + if (matchGroup(sourceEntity, clusterNodeEntity)) { + StreamSourceEntity fileEntity = + CommonBeanUtils.copyProperties(sourceEntity, StreamSourceEntity::new); + fileEntity.setSourceName(fileEntity.getSourceName() + "-" + + RandomStringUtils.randomAlphanumeric(10).toLowerCase(Locale.ROOT)); + fileEntity.setTemplateId(sourceEntity.getId()); + fileEntity.setAgentIp(agentIp); + fileEntity.setStatus(SourceStatus.TO_BE_ISSUED_ADD.getCode()); + // create new sub source task + sourceMapper.insert(fileEntity); + LOGGER.info("Transform new template task({}) for agent({}) in cluster({}).", + fileEntity.getId(), taskRequest.getAgentIp(), taskRequest.getClusterName()); + } + } + }); } /** diff --git a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java index 0fdae8b096..d27c018653 100644 --- a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java +++ b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/AgentServiceTest.java @@ -71,8 +71,6 @@ class AgentServiceTest extends ServiceBaseTest { @Autowired private AgentService agentService; @Autowired - private HeartbeatService heartbeatService; - @Autowired private InlongGroupEntityMapper groupMapper; @Autowired private InlongStreamServiceTest streamServiceTest; @@ -212,12 +210,12 @@ class AgentServiceTest extends ServiceBaseTest { TaskResult taskResult = agent.pullTask(); Assertions.assertTrue(taskResult.getCmdConfigs().isEmpty()); - Assertions.assertEquals(4, taskResult.getDataConfigs().size()); + Assertions.assertEquals(3, taskResult.getDataConfigs().size()); Assertions.assertEquals(3, taskResult.getDataConfigs().stream() .filter(dataConfig -> Integer.valueOf(dataConfig.getOp()) == ManagerOpEnum.ADD.getType()) .collect(Collectors.toSet()) .size()); - Assertions.assertEquals(1, taskResult.getDataConfigs().stream() + Assertions.assertEquals(0, taskResult.getDataConfigs().stream() .filter(dataConfig -> Integer.valueOf(dataConfig.getOp()) == ManagerOpEnum.FROZEN.getType()) .collect(Collectors.toSet()) .size());
