Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import org.apache.dolphinscheduler.api.enums.Status;
import org.apache.dolphinscheduler.api.exceptions.ServiceException;
import org.apache.dolphinscheduler.api.service.EnvironmentService;
import org.apache.dolphinscheduler.api.service.WorkerGroupService;
import org.apache.dolphinscheduler.api.utils.PageInfo;
import org.apache.dolphinscheduler.api.utils.Result;
import org.apache.dolphinscheduler.common.enums.AuthorizationType;
Expand All @@ -35,15 +36,18 @@
import org.apache.dolphinscheduler.dao.entity.EnvironmentWorkerGroupRelation;
import org.apache.dolphinscheduler.dao.entity.TaskDefinition;
import org.apache.dolphinscheduler.dao.entity.User;
import org.apache.dolphinscheduler.dao.entity.WorkerGroup;
import org.apache.dolphinscheduler.dao.mapper.EnvironmentMapper;
import org.apache.dolphinscheduler.dao.mapper.EnvironmentWorkerGroupRelationMapper;
import org.apache.dolphinscheduler.dao.repository.TaskDefinitionDao;
import org.apache.dolphinscheduler.dao.repository.WorkerGroupDao;

import org.apache.commons.collections4.CollectionUtils;
import org.apache.commons.collections4.SetUtils;
import org.apache.commons.lang3.StringUtils;

import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.Date;
import java.util.List;
Expand Down Expand Up @@ -82,6 +86,12 @@ public class EnvironmentServiceImpl extends BaseServiceImpl implements Environme
@Autowired
private TaskDefinitionDao taskDefinitionDao;

@Autowired
private WorkerGroupDao workerGroupDao;

@Autowired
private WorkerGroupService workerGroupService;

/**
* create environment
*
Expand Down Expand Up @@ -111,6 +121,9 @@ public Long createEnvironment(User loginUser,
throw new ServiceException(Status.ENVIRONMENT_NAME_EXISTS, name);
}

List<String> workerGroupList = parseWorkerGroupList(workerGroups);
Map<String, WorkerGroup> workerGroupMap = queryWorkerGroupMap(workerGroupList);

Environment env = new Environment();
env.setName(name);
env.setConfig(config);
Expand All @@ -121,25 +134,20 @@ public Long createEnvironment(User loginUser,
env.setCode(CodeGenerateUtils.genCode());

if (environmentMapper.insert(env) > 0) {
if (!StringUtils.isEmpty(workerGroups)) {
List<String> workerGroupList = JSONUtils.parseObject(workerGroups, new TypeReference<List<String>>() {
if (CollectionUtils.isNotEmpty(workerGroupList)) {
workerGroupList.forEach(workerGroup -> {
EnvironmentWorkerGroupRelation relation = new EnvironmentWorkerGroupRelation();
relation.setEnvironmentCode(env.getCode());
relation.setWorkerGroupId(workerGroupMap.get(workerGroup).getId());
relation.setWorkerGroup(workerGroup);
relation.setOperator(loginUser.getId());
relation.setCreateTime(new Date());
relation.setUpdateTime(new Date());
relationMapper.insert(relation);
log.info(
"Environment-WorkerGroup relation create complete, environmentCode:{}, workerGroupId:{}.",
env.getCode(), relation.getWorkerGroupId());
});
if (CollectionUtils.isNotEmpty(workerGroupList)) {
workerGroupList.stream().forEach(workerGroup -> {
if (!StringUtils.isEmpty(workerGroup)) {
EnvironmentWorkerGroupRelation relation = new EnvironmentWorkerGroupRelation();
relation.setEnvironmentCode(env.getCode());
relation.setWorkerGroup(workerGroup);
relation.setOperator(loginUser.getId());
relation.setCreateTime(new Date());
relation.setUpdateTime(new Date());
relationMapper.insert(relation);
log.info(
"Environment-WorkerGroup relation create complete, environmentName:{}, workerGroup:{}.",
env.getName(), relation.getWorkerGroup());
}
});
}
}
return env.getCode();
}
Expand Down Expand Up @@ -329,13 +337,7 @@ public Environment updateEnvironmentByCode(User loginUser, Long code, String nam
throw new ServiceException(Status.ENVIRONMENT_NAME_EXISTS, name);
}

Set<String> workerGroupSet;
if (!StringUtils.isEmpty(workerGroups)) {
workerGroupSet = JSONUtils.parseObject(workerGroups, new TypeReference<Set<String>>() {
});
} else {
workerGroupSet = new TreeSet<>();
}
Set<String> workerGroupSet = new TreeSet<>(parseWorkerGroupList(workerGroups));

Set<String> existWorkerGroupSet = relationMapper
.queryByEnvironmentCode(code)
Expand All @@ -345,6 +347,7 @@ public Environment updateEnvironmentByCode(User loginUser, Long code, String nam

Set<String> deleteWorkerGroupSet = SetUtils.difference(existWorkerGroupSet, workerGroupSet).toSet();
Set<String> addWorkerGroupSet = SetUtils.difference(workerGroupSet, existWorkerGroupSet).toSet();
Map<String, WorkerGroup> workerGroupMap = queryWorkerGroupMap(addWorkerGroupSet);

// verify whether the relation of this environment and worker groups can be adjusted
checkUsedEnvironmentWorkerGroupRelation(deleteWorkerGroupSet, name, code);
Expand Down Expand Up @@ -374,6 +377,7 @@ public Environment updateEnvironmentByCode(User loginUser, Long code, String nam
if (StringUtils.isNotEmpty(key)) {
EnvironmentWorkerGroupRelation relation = new EnvironmentWorkerGroupRelation();
relation.setEnvironmentCode(code);
relation.setWorkerGroupId(workerGroupMap.get(key).getId());
relation.setWorkerGroup(key);
relation.setUpdateTime(new Date());
relation.setCreateTime(new Date());
Expand Down Expand Up @@ -418,6 +422,49 @@ private void checkUsedEnvironmentWorkerGroupRelation(Set<String> deleteKeySet,
}
}

private List<String> parseWorkerGroupList(String workerGroups) {
if (StringUtils.isEmpty(workerGroups)) {
return Collections.emptyList();
}
List<String> workerGroupList = JSONUtils.parseObject(workerGroups, new TypeReference<List<String>>() {
});
if (CollectionUtils.isEmpty(workerGroupList)) {
return Collections.emptyList();
}
return workerGroupList.stream()
.filter(StringUtils::isNotEmpty)
.collect(Collectors.toList());
}

private Map<String, WorkerGroup> queryWorkerGroupMap(Collection<String> workerGroupNames) {
if (CollectionUtils.isEmpty(workerGroupNames)) {
return Collections.emptyMap();
}
Set<String> nonEmptyWorkerGroupNames = workerGroupNames.stream()
.filter(StringUtils::isNotEmpty)
.collect(Collectors.toCollection(TreeSet::new));
if (CollectionUtils.isEmpty(nonEmptyWorkerGroupNames)) {
return Collections.emptyMap();
}

List<WorkerGroup> workerGroups = workerGroupDao.queryWorkerGroupByNames(nonEmptyWorkerGroupNames);
Map<String, WorkerGroup> workerGroupMap = CollectionUtils.emptyIfNull(workerGroups).stream()
.collect(Collectors.toMap(WorkerGroup::getName, workerGroup -> workerGroup));
if (!workerGroupMap.keySet().containsAll(nonEmptyWorkerGroupNames)) {
// Configuration-defined groups are identified by name and have no database id.
workerGroupService.getConfigWorkerGroupPageDetail().stream()
.filter(workerGroup -> nonEmptyWorkerGroupNames.contains(workerGroup.getName()))
.forEach(workerGroup -> workerGroupMap.putIfAbsent(workerGroup.getName(), workerGroup));
}
Set<String> notExistWorkerGroups =
SetUtils.difference(nonEmptyWorkerGroupNames, workerGroupMap.keySet()).toSet();
if (CollectionUtils.isNotEmpty(notExistWorkerGroups)) {
throw new ServiceException(Status.WORKER_GROUP_NOT_EXIST,
String.join(",", new TreeSet<>(notExistWorkerGroups)));
}
return workerGroupMap;
}

protected void checkParams(String name, String config, String workerGroups) {
if (StringUtils.isEmpty(name)) {
throw new ServiceException(Status.ENVIRONMENT_NAME_IS_NULL);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -182,7 +182,10 @@ private void checkWorkerGroupDependencies(WorkerGroup workerGroup) {
// check if the worker group has any dependent environments
List<EnvironmentWorkerGroupRelation> environmentWorkerGroupRelations =
environmentWorkerGroupRelationMapper.selectList(new QueryWrapper<EnvironmentWorkerGroupRelation>()
.lambda().eq(EnvironmentWorkerGroupRelation::getWorkerGroup, workerGroup.getName()));
.lambda()
.eq(EnvironmentWorkerGroupRelation::getWorkerGroupId, workerGroup.getId())
.or()
.eq(EnvironmentWorkerGroupRelation::getWorkerGroup, workerGroup.getName()));

if (CollectionUtils.isNotEmpty(environmentWorkerGroupRelations)) {
throw new ServiceException(Status.WORKER_GROUP_DEPENDENT_ENVIRONMENT_EXISTS,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import org.apache.dolphinscheduler.common.enums.AuthorizationType;
import org.apache.dolphinscheduler.common.enums.UserType;
import org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus;
import org.apache.dolphinscheduler.dao.entity.EnvironmentWorkerGroupRelation;
import org.apache.dolphinscheduler.dao.entity.User;
import org.apache.dolphinscheduler.dao.entity.WorkerGroup;
import org.apache.dolphinscheduler.dao.mapper.EnvironmentWorkerGroupRelationMapper;
Expand All @@ -43,14 +44,18 @@
import org.apache.dolphinscheduler.registry.api.RegistryClient;
import org.apache.dolphinscheduler.registry.api.enums.RegistryNodeType;

import org.apache.ibatis.builder.MapperBuilderAssistant;

import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;

import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
Expand All @@ -63,6 +68,10 @@
import org.slf4j.LoggerFactory;
import org.springframework.dao.DuplicateKeyException;

import com.baomidou.mybatisplus.core.MybatisConfiguration;
import com.baomidou.mybatisplus.core.conditions.Wrapper;
import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;

@ExtendWith(MockitoExtension.class)
@MockitoSettings(strictness = Strictness.LENIENT)
public class WorkerGroupServiceTest {
Expand Down Expand Up @@ -97,6 +106,14 @@ public class WorkerGroupServiceTest {

private final String GROUP_NAME = "testWorkerGroup";

@BeforeEach
public void setUp() {
if (TableInfoHelper.getTableInfo(EnvironmentWorkerGroupRelation.class) == null) {
TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new MybatisConfiguration(), ""),
EnvironmentWorkerGroupRelation.class);
}
}

private User getLoginUser() {
User loginUser = new User();
loginUser.setUserType(UserType.GENERAL_USER);
Expand Down Expand Up @@ -260,6 +277,66 @@ public void giveValidParams_whenDeleteWorkerGroupById_expectSuccess() {
assertDoesNotThrow(() -> workerGroupService.deleteWorkerGroupById(loginUser, 1));
}

@Test
public void giveEnvironmentRelationMatchedByWorkerGroupId_whenDeleteWorkerGroupById_expectEnvironmentDependency() {
User loginUser = getLoginUser();
when(resourcePermissionCheckService.operationPermissionCheck(AuthorizationType.WORKER_GROUP, 1,
WORKER_GROUP_DELETE, baseServiceLogger)).thenReturn(true);
when(resourcePermissionCheckService.resourcePermissionCheck(AuthorizationType.WORKER_GROUP, null, 1,
baseServiceLogger)).thenReturn(true);
WorkerGroup workerGroup = getWorkerGroup(1);
when(workerGroupDao.queryById(1)).thenReturn(workerGroup);
when(workflowInstanceDao.queryByWorkerGroupNameAndStatus(workerGroup.getName(),
WorkflowExecutionStatus.NOT_TERMINAL_STATES)).thenReturn(null);
when(taskDefinitionDao.queryByWorkerGroup(Mockito.any())).thenReturn(null);
when(scheduleDao.queryScheduleByWorkerGroup(Mockito.any())).thenReturn(null);

EnvironmentWorkerGroupRelation relation = new EnvironmentWorkerGroupRelation();
relation.setWorkerGroupId(workerGroup.getId());
relation.setWorkerGroup("stale-" + workerGroup.getName());
when(environmentWorkerGroupRelationMapper.selectList(Mockito.any())).thenAnswer(invocation -> {
Wrapper<EnvironmentWorkerGroupRelation> wrapper = invocation.getArgument(0);
String sqlSegment = wrapper.getSqlSegment();
if (sqlSegment.contains("worker_group_id") && sqlSegment.contains("worker_group")) {
return Collections.singletonList(relation);
}
return Collections.emptyList();
});

assertThrowsServiceException(Status.WORKER_GROUP_DEPENDENT_ENVIRONMENT_EXISTS,
() -> workerGroupService.deleteWorkerGroupById(loginUser, 1));
}

@Test
public void giveEnvironmentRelationMatchedOnlyByWorkerGroupNameWithDifferentWorkerGroupId_whenDeleteWorkerGroupById_expectEnvironmentDependency() {
User loginUser = getLoginUser();
when(resourcePermissionCheckService.operationPermissionCheck(AuthorizationType.WORKER_GROUP, 1,
WORKER_GROUP_DELETE, baseServiceLogger)).thenReturn(true);
when(resourcePermissionCheckService.resourcePermissionCheck(AuthorizationType.WORKER_GROUP, null, 1,
baseServiceLogger)).thenReturn(true);
WorkerGroup workerGroup = getWorkerGroup(1);
when(workerGroupDao.queryById(1)).thenReturn(workerGroup);
when(workflowInstanceDao.queryByWorkerGroupNameAndStatus(workerGroup.getName(),
WorkflowExecutionStatus.NOT_TERMINAL_STATES)).thenReturn(null);
when(taskDefinitionDao.queryByWorkerGroup(Mockito.any())).thenReturn(null);
when(scheduleDao.queryScheduleByWorkerGroup(Mockito.any())).thenReturn(null);

EnvironmentWorkerGroupRelation relation = new EnvironmentWorkerGroupRelation();
relation.setWorkerGroupId(2);
relation.setWorkerGroup(workerGroup.getName());
when(environmentWorkerGroupRelationMapper.selectList(Mockito.any())).thenAnswer(invocation -> {
Wrapper<EnvironmentWorkerGroupRelation> wrapper = invocation.getArgument(0);
String sqlSegment = wrapper.getSqlSegment();
Assertions.assertTrue(sqlSegment.contains("worker_group_id"));
Assertions.assertTrue(sqlSegment.contains("worker_group =") || sqlSegment.contains("worker_group="));
Assertions.assertFalse(sqlSegment.contains("IS NULL"));
return Collections.singletonList(relation);
});

assertThrowsServiceException(Status.WORKER_GROUP_DEPENDENT_ENVIRONMENT_EXISTS,
() -> workerGroupService.deleteWorkerGroupById(loginUser, 1));
}

@Test
public void testQueryAllGroupWithDefault() {
List<String> workerGroups = workerGroupService.queryAllGroup(getLoginUser());
Expand Down
Loading
Loading