diff --git a/.github/workflows/cluster-test/mysql_with_mysql_registry/dolphinscheduler_env.sh b/.github/workflows/cluster-test/mysql_with_mysql_registry/dolphinscheduler_env.sh
index 58937e740c..8536eb0905 100755
--- a/.github/workflows/cluster-test/mysql_with_mysql_registry/dolphinscheduler_env.sh
+++ b/.github/workflows/cluster-test/mysql_with_mysql_registry/dolphinscheduler_env.sh
@@ -28,7 +28,6 @@ export SPRING_DATASOURCE_PASSWORD=123456
# DolphinScheduler server related configuration
export SPRING_CACHE_TYPE=${SPRING_CACHE_TYPE:-none}
export SPRING_JACKSON_TIME_ZONE=${SPRING_JACKSON_TIME_ZONE:-UTC}
-export MASTER_FETCH_COMMAND_NUM=${MASTER_FETCH_COMMAND_NUM:-10}
# Registry center configuration, determines the type and link of the registry center
export REGISTRY_TYPE=${REGISTRY_TYPE:-jdbc}
diff --git a/.github/workflows/cluster-test/mysql_with_zookeeper_registry/dolphinscheduler_env.sh b/.github/workflows/cluster-test/mysql_with_zookeeper_registry/dolphinscheduler_env.sh
index 671c70a5bb..f64e59b768 100755
--- a/.github/workflows/cluster-test/mysql_with_zookeeper_registry/dolphinscheduler_env.sh
+++ b/.github/workflows/cluster-test/mysql_with_zookeeper_registry/dolphinscheduler_env.sh
@@ -28,7 +28,6 @@ export SPRING_DATASOURCE_PASSWORD=123456
# DolphinScheduler server related configuration
export SPRING_CACHE_TYPE=${SPRING_CACHE_TYPE:-none}
export SPRING_JACKSON_TIME_ZONE=${SPRING_JACKSON_TIME_ZONE:-UTC}
-export MASTER_FETCH_COMMAND_NUM=${MASTER_FETCH_COMMAND_NUM:-10}
# Registry center configuration, determines the type and link of the registry center
export REGISTRY_TYPE=${REGISTRY_TYPE:-zookeeper}
diff --git a/.github/workflows/cluster-test/postgresql_with_postgresql_registry/dolphinscheduler_env.sh b/.github/workflows/cluster-test/postgresql_with_postgresql_registry/dolphinscheduler_env.sh
index e7fd1b7204..29f8570319 100644
--- a/.github/workflows/cluster-test/postgresql_with_postgresql_registry/dolphinscheduler_env.sh
+++ b/.github/workflows/cluster-test/postgresql_with_postgresql_registry/dolphinscheduler_env.sh
@@ -28,7 +28,6 @@ export SPRING_DATASOURCE_PASSWORD=postgres
# DolphinScheduler server related configuration
export SPRING_CACHE_TYPE=${SPRING_CACHE_TYPE:-none}
export SPRING_JACKSON_TIME_ZONE=${SPRING_JACKSON_TIME_ZONE:-UTC}
-export MASTER_FETCH_COMMAND_NUM=${MASTER_FETCH_COMMAND_NUM:-10}
# Registry center configuration, determines the type and link of the registry center
export REGISTRY_TYPE=jdbc
diff --git a/.github/workflows/cluster-test/postgresql_with_zookeeper_registry/dolphinscheduler_env.sh b/.github/workflows/cluster-test/postgresql_with_zookeeper_registry/dolphinscheduler_env.sh
index 1dbd63254e..6851716058 100644
--- a/.github/workflows/cluster-test/postgresql_with_zookeeper_registry/dolphinscheduler_env.sh
+++ b/.github/workflows/cluster-test/postgresql_with_zookeeper_registry/dolphinscheduler_env.sh
@@ -28,7 +28,6 @@ export SPRING_DATASOURCE_PASSWORD=postgres
# DolphinScheduler server related configuration
export SPRING_CACHE_TYPE=${SPRING_CACHE_TYPE:-none}
export SPRING_JACKSON_TIME_ZONE=${SPRING_JACKSON_TIME_ZONE:-UTC}
-export MASTER_FETCH_COMMAND_NUM=${MASTER_FETCH_COMMAND_NUM:-10}
# Registry center configuration, determines the type and link of the registry center
export REGISTRY_TYPE=${REGISTRY_TYPE:-zookeeper}
diff --git a/docs/docs/en/architecture/configuration.md b/docs/docs/en/architecture/configuration.md
index 13d8932943..fe0b7851ba 100644
--- a/docs/docs/en/architecture/configuration.md
+++ b/docs/docs/en/architecture/configuration.md
@@ -286,7 +286,6 @@ Location: `master-server/conf/application.yaml`
| Parameters | Default value | Description |
|-----------------------------------------------------------------------------|---------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
| master.listen-port | 5678 | master listen port |
-| master.fetch-command-num | 10 | the number of commands fetched by master |
| master.pre-exec-threads | 10 | master prepare execute thread number to limit handle commands in parallel |
| master.exec-threads | 100 | master execute thread number to limit process instances in parallel |
| master.dispatch-task-number | 3 | master dispatch task number per batch |
@@ -305,6 +304,9 @@ Location: `master-server/conf/application.yaml`
| master.registry-disconnect-strategy.strategy | stop | Used when the master disconnect from registry, default value: stop. Optional values include stop, waiting |
| master.registry-disconnect-strategy.max-waiting-time | 100s | Used when the master disconnect from registry, and the disconnect strategy is waiting, this config means the master will waiting to reconnect to registry in given times, and after the waiting times, if the master still cannot connect to registry, will stop itself, if the value is 0s, the Master will wait infinitely |
| master.worker-group-refresh-interval | 10s | The interval to refresh worker group from db to memory |
+| master.command-fetch-strategy.type | ID_SLOT_BASED | The command fetch strategy, only support `ID_SLOT_BASED` |
+| master.command-fetch-strategy.config.id-step | 1 | The id auto incremental step of t_ds_command in db |
+| master.command-fetch-strategy.config.fetch-size | 10 | The number of commands fetched by master |
### Worker Server related configuration
diff --git a/docs/docs/en/guide/installation/pseudo-cluster.md b/docs/docs/en/guide/installation/pseudo-cluster.md
index e63436f203..7a3b43b00e 100644
--- a/docs/docs/en/guide/installation/pseudo-cluster.md
+++ b/docs/docs/en/guide/installation/pseudo-cluster.md
@@ -123,7 +123,6 @@ export SPRING_DATASOURCE_PASSWORD={password}
# DolphinScheduler server related configuration
export SPRING_CACHE_TYPE=${SPRING_CACHE_TYPE:-none}
export SPRING_JACKSON_TIME_ZONE=${SPRING_JACKSON_TIME_ZONE:-UTC}
-export MASTER_FETCH_COMMAND_NUM=${MASTER_FETCH_COMMAND_NUM:-10}
# Registry center configuration, determines the type and link of the registry center
export REGISTRY_TYPE=${REGISTRY_TYPE:-zookeeper}
diff --git a/docs/docs/zh/architecture/configuration.md b/docs/docs/zh/architecture/configuration.md
index 08fded19e0..d8d1d42d1e 100644
--- a/docs/docs/zh/architecture/configuration.md
+++ b/docs/docs/zh/architecture/configuration.md
@@ -281,29 +281,30 @@ common.properties配置文件目前主要是配置hadoop/s3/yarn/applicationId
位置:`master-server/conf/application.yaml`
-| 参数 | 默认值 | 描述 |
-|-----------------------------------------------------------------------------|--------------|-----------------------------------------------------------------------------------------|
-| master.listen-port | 5678 | master监听端口 |
-| master.fetch-command-num | 10 | master拉取command数量 |
-| master.pre-exec-threads | 10 | master准备执行任务的数量,用于限制并行的command |
-| master.exec-threads | 100 | master工作线程数量,用于限制并行的流程实例数量 |
-| master.dispatch-task-number | 3 | master每个批次的派发任务数量 |
-| master.host-selector | lower_weight | master host选择器,用于选择合适的worker执行任务,可选值: random, round_robin, lower_weight |
-| master.max-heartbeat-interval | 10s | master最大心跳间隔 |
-| master.task-commit-retry-times | 5 | 任务重试次数 |
-| master.task-commit-interval | 1000 | 任务提交间隔,单位为毫秒 |
-| master.state-wheel-interval | 5 | 轮询检查状态时间 |
-| master.server-load-protection.enabled | true | 是否开启系统保护策略 |
-| master.server-load-protection.max-system-cpu-usage-percentage-thresholds | 0.7 | master最大系统cpu使用值,只有当前系统cpu使用值低于最大系统cpu使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的操作系统CPU |
-| master.server-load-protection.max-jvm-cpu-usage-percentage-thresholds | 0.7 | master最大JVM cpu使用值,只有当前JVM cpu使用值低于最大JVM cpu使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的JVM CPU |
-| master.server-load-protection.max-system-memory-usage-percentage-thresholds | 0.7 | master最大系统 内存使用值,只有当前系统内存使用值低于最大系统内存使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的操作系统内存 |
-| master.server-load-protection.max-disk-usage-percentage-thresholds | 0.7 | master最大系统磁盘使用值,只有当前系统磁盘使用值低于最大系统磁盘使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的操作系统磁盘空间 |
-| master.failover-interval | 10 | failover间隔,单位为分钟 |
-| master.kill-application-when-task-failover | true | 当任务实例failover时,是否kill掉yarn或k8s application |
-| master.registry-disconnect-strategy.strategy | stop | 当Master与注册中心失联之后采取的策略, 默认值是: stop. 可选值包括: stop, waiting |
-| master.registry-disconnect-strategy.max-waiting-time | 100s | 当Master与注册中心失联之后重连时间, 之后当strategy为waiting时,该值生效。 该值表示当Master与注册中心失联时会在给定时间之内进行重连, |
-| 在给定时间之内重连失败将会停止自己,在重连时,Master会丢弃目前正在执行的工作流,值为0表示会无限期等待 |
-| master.master.worker-group-refresh-interval | 10s | 定期将workerGroup从数据库中同步到内存的时间间隔 |
+| 参数 | 默认值 | 描述 |
+|-----------------------------------------------------------------------------|---------------|------------------------------------------------------------------------------------------------------------------------------------------|
+| master.listen-port | 5678 | master监听端口 |
+| master.pre-exec-threads | 10 | master准备执行任务的数量,用于限制并行的command |
+| master.exec-threads | 100 | master工作线程数量,用于限制并行的流程实例数量 |
+| master.dispatch-task-number | 3 | master每个批次的派发任务数量 |
+| master.host-selector | lower_weight | master host选择器,用于选择合适的worker执行任务,可选值: random, round_robin, lower_weight |
+| master.max-heartbeat-interval | 10s | master最大心跳间隔 |
+| master.task-commit-retry-times | 5 | 任务重试次数 |
+| master.task-commit-interval | 1000 | 任务提交间隔,单位为毫秒 |
+| master.state-wheel-interval | 5 | 轮询检查状态时间 |
+| master.server-load-protection.enabled | true | 是否开启系统保护策略 |
+| master.server-load-protection.max-system-cpu-usage-percentage-thresholds | 0.7 | master最大系统cpu使用值,只有当前系统cpu使用值低于最大系统cpu使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的操作系统CPU |
+| master.server-load-protection.max-jvm-cpu-usage-percentage-thresholds | 0.7 | master最大JVM cpu使用值,只有当前JVM cpu使用值低于最大JVM cpu使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的JVM CPU |
+| master.server-load-protection.max-system-memory-usage-percentage-thresholds | 0.7 | master最大系统 内存使用值,只有当前系统内存使用值低于最大系统内存使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的操作系统内存 |
+| master.server-load-protection.max-disk-usage-percentage-thresholds | 0.7 | master最大系统磁盘使用值,只有当前系统磁盘使用值低于最大系统磁盘使用值,master服务才能调度任务. 默认值为0.7: 会使用70%的操作系统磁盘空间 |
+| master.failover-interval | 10 | failover间隔,单位为分钟 |
+| master.kill-application-when-task-failover | true | 当任务实例failover时,是否kill掉yarn或k8s application |
+| master.registry-disconnect-strategy.strategy | stop | 当Master与注册中心失联之后采取的策略, 默认值是: stop. 可选值包括: stop, waiting |
+| master.registry-disconnect-strategy.max-waiting-time | 100s | 当Master与注册中心失联之后重连时间, 之后当strategy为waiting时,该值生效。 该值表示当Master与注册中心失联时会在给定时间之内进行重连, 在给定时间之内重连失败将会停止自己,在重连时,Master会丢弃目前正在执行的工作流,值为0表示会无限期等待 |
+| master.master.worker-group-refresh-interval | 10s | 定期将workerGroup从数据库中同步到内存的时间间隔 |
+| master.command-fetch-strategy.type | ID_SLOT_BASED | Command拉取策略, 目前仅支持 `ID_SLOT_BASED` |
+| master.command-fetch-strategy.config.id-step | 1 | 数据库中t_ds_command的id自增步长 |
+| master.command-fetch-strategy.config.fetch-size | 10 | master拉取command数量 |
## Worker Server相关配置
diff --git a/docs/docs/zh/guide/installation/pseudo-cluster.md b/docs/docs/zh/guide/installation/pseudo-cluster.md
index 13479e0d9e..a199167e04 100644
--- a/docs/docs/zh/guide/installation/pseudo-cluster.md
+++ b/docs/docs/zh/guide/installation/pseudo-cluster.md
@@ -118,7 +118,6 @@ export SPRING_DATASOURCE_PASSWORD={password}
# DolphinScheduler server related configuration
export SPRING_CACHE_TYPE=${SPRING_CACHE_TYPE:-none}
export SPRING_JACKSON_TIME_ZONE=${SPRING_JACKSON_TIME_ZONE:-UTC}
-export MASTER_FETCH_COMMAND_NUM=${MASTER_FETCH_COMMAND_NUM:-10}
# Registry center configuration, determines the type and link of the registry center
export REGISTRY_TYPE=${REGISTRY_TYPE:-zookeeper}
diff --git a/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/assembly/dolphinscheduler-alert-server.xml b/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/assembly/dolphinscheduler-alert-server.xml
index 6f39096221..bf28193b3f 100644
--- a/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/assembly/dolphinscheduler-alert-server.xml
+++ b/dolphinscheduler-alert/dolphinscheduler-alert-server/src/main/assembly/dolphinscheduler-alert-server.xml
@@ -52,6 +52,7 @@
${basedir}/../../dolphinscheduler-common/src/main/resources**/*.properties
+ **/*.yamlconf
diff --git a/dolphinscheduler-api/src/main/assembly/dolphinscheduler-api-server.xml b/dolphinscheduler-api/src/main/assembly/dolphinscheduler-api-server.xml
index 77b2f54c64..c9fdab4d51 100644
--- a/dolphinscheduler-api/src/main/assembly/dolphinscheduler-api-server.xml
+++ b/dolphinscheduler-api/src/main/assembly/dolphinscheduler-api-server.xml
@@ -53,6 +53,7 @@
${basedir}/../dolphinscheduler-common/src/main/resources**/*.properties
+ **/*.yamlconf
diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProjectWorkerGroupRelationServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProjectWorkerGroupRelationServiceImpl.java
index 15cbcbb7fa..31dc113dbb 100644
--- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProjectWorkerGroupRelationServiceImpl.java
+++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProjectWorkerGroupRelationServiceImpl.java
@@ -26,6 +26,7 @@ import org.apache.dolphinscheduler.common.constants.Constants;
import org.apache.dolphinscheduler.dao.entity.Project;
import org.apache.dolphinscheduler.dao.entity.ProjectWorkerGroup;
import org.apache.dolphinscheduler.dao.entity.User;
+import org.apache.dolphinscheduler.dao.entity.WorkerGroup;
import org.apache.dolphinscheduler.dao.mapper.ProjectMapper;
import org.apache.dolphinscheduler.dao.mapper.ProjectWorkerGroupMapper;
import org.apache.dolphinscheduler.dao.mapper.ScheduleMapper;
@@ -38,6 +39,7 @@ import org.apache.commons.lang3.StringUtils;
import java.util.Date;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -113,23 +115,25 @@ public class ProjectWorkerGroupRelationServiceImpl extends BaseServiceImpl
}
Set workerGroupNames =
- workerGroupMapper.queryAllWorkerGroup().stream().map(item -> item.getName()).collect(
+ workerGroupMapper.queryAllWorkerGroup().stream().map(WorkerGroup::getName).collect(
Collectors.toSet());
workerGroupNames.add(Constants.DEFAULT_WORKER_GROUP);
- Set assignedWorkerGroupNames = workerGroups.stream().collect(Collectors.toSet());
+ Set assignedWorkerGroupNames = new HashSet<>(workerGroups);
Set difference = SetUtils.difference(assignedWorkerGroupNames, workerGroupNames);
- if (difference.size() > 0) {
+ if (!difference.isEmpty()) {
putMsg(result, Status.WORKER_GROUP_NOT_EXIST, difference.toString());
return result;
}
Set projectWorkerGroupNames = projectWorkerGroupMapper.selectList(new QueryWrapper()
.lambda()
- .eq(ProjectWorkerGroup::getProjectCode, projectCode)).stream().map(item -> item.getWorkerGroup())
+ .eq(ProjectWorkerGroup::getProjectCode, projectCode))
+ .stream()
+ .map(ProjectWorkerGroup::getWorkerGroup)
.collect(Collectors.toSet());
difference = SetUtils.difference(projectWorkerGroupNames, assignedWorkerGroupNames);
diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectPreferenceServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectPreferenceServiceTest.java
index 530c15d48e..7a74a7c265 100644
--- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectPreferenceServiceTest.java
+++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectPreferenceServiceTest.java
@@ -60,28 +60,65 @@ public class ProjectPreferenceServiceTest {
public void testUpdateProjectPreference() {
User loginUser = getGeneralUser();
+ // no permission
+ Mockito.when(projectService.hasProjectAndWritePerm(Mockito.any(), Mockito.any(), Mockito.any(Result.class)))
+ .thenReturn(false);
+ Result result = projectPreferenceService.updateProjectPreference(loginUser, projectCode, "value");
+ Assertions.assertNull(result.getCode());
+ Assertions.assertNull(result.getData());
+ Assertions.assertNull(result.getMsg());
+
+ // when preference exists in project
+ Mockito.when(projectPreferenceMapper.selectOne(Mockito.any())).thenReturn(null);
Mockito.when(projectMapper.queryByCode(projectCode)).thenReturn(getProject(projectCode));
+
+ // success
Mockito.when(projectService.hasProjectAndWritePerm(Mockito.any(), Mockito.any(), Mockito.any(Result.class)))
.thenReturn(true);
- Mockito.when(projectPreferenceMapper.selectOne(Mockito.any())).thenReturn(null);
Mockito.when(projectPreferenceMapper.insert(Mockito.any())).thenReturn(1);
- Result result = projectPreferenceService.updateProjectPreference(loginUser, projectCode, "value");
+ result = projectPreferenceService.updateProjectPreference(loginUser, projectCode, "value");
Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode());
+
+ // database operatation fail
+ Mockito.when(projectPreferenceMapper.insert(Mockito.any())).thenReturn(-1);
+ result = projectPreferenceService.updateProjectPreference(loginUser, projectCode, "value");
+ Assertions.assertEquals(Status.CREATE_PROJECT_PREFERENCE_ERROR.getCode(), result.getCode());
+
+ // when preference exists in project
+ Mockito.when(projectPreferenceMapper.selectOne(Mockito.any())).thenReturn(getProjectPreference());
+
+ // success
+ Mockito.when(projectPreferenceMapper.updateById(Mockito.any())).thenReturn(1);
+ result = projectPreferenceService.updateProjectPreference(loginUser, projectCode, "value");
+ Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode());
+
+ // database operation fail
+ Mockito.when(projectPreferenceMapper.updateById(Mockito.any())).thenReturn(-1);
+ result = projectPreferenceService.updateProjectPreference(loginUser, projectCode, "value");
+ Assertions.assertEquals(Status.UPDATE_PROJECT_PREFERENCE_ERROR.getCode(), result.getCode());
}
@Test
public void testQueryProjectPreferenceByProjectCode() {
User loginUser = getGeneralUser();
+ // no permission
+ Mockito.when(projectService.hasProjectAndWritePerm(Mockito.any(), Mockito.any(), Mockito.any(Result.class)))
+ .thenReturn(false);
+ Result result = projectPreferenceService.queryProjectPreferenceByProjectCode(loginUser, projectCode);
+ Assertions.assertNull(result.getCode());
+ Assertions.assertNull(result.getData());
+ Assertions.assertNull(result.getMsg());
+
// PROJECT_PARAMETER_NOT_EXISTS
Mockito.when(projectMapper.queryByCode(projectCode)).thenReturn(getProject(projectCode));
Mockito.when(projectService.hasProjectAndPerm(Mockito.any(), Mockito.any(), Mockito.any(Result.class),
Mockito.any())).thenReturn(true);
Mockito.when(projectPreferenceMapper.selectOne(Mockito.any())).thenReturn(null);
- Result result = projectPreferenceService.queryProjectPreferenceByProjectCode(loginUser, projectCode);
+ result = projectPreferenceService.queryProjectPreferenceByProjectCode(loginUser, projectCode);
Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode());
// SUCCESS
@@ -94,14 +131,29 @@ public class ProjectPreferenceServiceTest {
public void testEnableProjectPreference() {
User loginUser = getGeneralUser();
+ // no permission
+ Mockito.when(projectService.hasProjectAndWritePerm(Mockito.any(), Mockito.any(), Mockito.any(Result.class)))
+ .thenReturn(false);
+ Result result = projectPreferenceService.enableProjectPreference(loginUser, projectCode, 1);
+ Assertions.assertNull(result.getCode());
+ Assertions.assertNull(result.getData());
+ Assertions.assertNull(result.getMsg());
+
Mockito.when(projectMapper.queryByCode(projectCode)).thenReturn(getProject(projectCode));
Mockito.when(projectService.hasProjectAndWritePerm(Mockito.any(), Mockito.any(), Mockito.any(Result.class)))
.thenReturn(true);
+ // success
Mockito.when(projectPreferenceMapper.selectOne(Mockito.any())).thenReturn(getProjectPreference());
- Result result = projectPreferenceService.enableProjectPreference(loginUser, projectCode, 1);
+ Mockito.when(projectPreferenceMapper.updateById(Mockito.any())).thenReturn(1);
+ result = projectPreferenceService.enableProjectPreference(loginUser, projectCode, 2);
Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode());
+ // db operation fail
+ Mockito.when(projectPreferenceMapper.selectOne(Mockito.any())).thenReturn(getProjectPreference());
+ Mockito.when(projectPreferenceMapper.updateById(Mockito.any())).thenReturn(-1);
+ result = projectPreferenceService.enableProjectPreference(loginUser, projectCode, 2);
+ Assertions.assertEquals(Status.UPDATE_PROJECT_PREFERENCE_STATE_ERROR.getCode(), result.getCode());
}
private User getGeneralUser() {
diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectWorkerGroupRelationServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectWorkerGroupRelationServiceTest.java
index a8e8b8beb2..6f6b7ca2d6 100644
--- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectWorkerGroupRelationServiceTest.java
+++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProjectWorkerGroupRelationServiceTest.java
@@ -17,13 +17,17 @@
package org.apache.dolphinscheduler.api.service;
+import static org.apache.dolphinscheduler.api.utils.ServiceTestUtil.getAdminUser;
+import static org.apache.dolphinscheduler.api.utils.ServiceTestUtil.getGeneralUser;
+
+import org.apache.dolphinscheduler.api.AssertionsHelper;
import org.apache.dolphinscheduler.api.enums.Status;
import org.apache.dolphinscheduler.api.service.impl.ProjectWorkerGroupRelationServiceImpl;
import org.apache.dolphinscheduler.api.utils.Result;
import org.apache.dolphinscheduler.common.constants.Constants;
-import org.apache.dolphinscheduler.common.enums.UserType;
import org.apache.dolphinscheduler.dao.entity.Project;
import org.apache.dolphinscheduler.dao.entity.ProjectWorkerGroup;
+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.ProjectMapper;
@@ -33,6 +37,7 @@ import org.apache.dolphinscheduler.dao.mapper.TaskDefinitionMapper;
import org.apache.dolphinscheduler.dao.mapper.WorkerGroupMapper;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
import java.util.Map;
@@ -77,27 +82,87 @@ public class ProjectWorkerGroupRelationServiceTest {
@Test
public void testAssignWorkerGroupsToProject() {
+ User generalUser = getGeneralUser();
User loginUser = getAdminUser();
- Mockito.when(projectMapper.queryByCode(projectCode)).thenReturn(null);
- Result result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ // no permission
+ Result result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(generalUser, projectCode,
+ getWorkerGroups());
+ Assertions.assertEquals(Status.USER_NO_OPERATION_PERM.getCode(), result.getCode());
+
+ // project code is null
+ result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, null,
getWorkerGroups());
Assertions.assertEquals(Status.PROJECT_NOT_EXIST.getCode(), result.getCode());
+ // worker group is empty
+ result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ Collections.emptyList());
+ Assertions.assertEquals(Status.WORKER_GROUP_TO_PROJECT_IS_EMPTY.getCode(), result.getCode());
+
+ // project not exists
+ Mockito.when(projectMapper.queryByCode(projectCode)).thenReturn(null);
+ result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ getWorkerGroups());
+ Assertions.assertEquals(Status.PROJECT_NOT_EXIST.getCode(), result.getCode());
+
+ // worker group not exists
WorkerGroup workerGroup = new WorkerGroup();
workerGroup.setName("test");
Mockito.when(projectMapper.queryByCode(Mockito.anyLong())).thenReturn(getProject());
- Mockito.when(workerGroupMapper.queryAllWorkerGroup()).thenReturn(Lists.newArrayList(workerGroup));
+ Mockito.when(workerGroupMapper.queryAllWorkerGroup()).thenReturn(Collections.singletonList(workerGroup));
+ result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ getDiffWorkerGroups());
+ Assertions.assertEquals(Status.WORKER_GROUP_NOT_EXIST.getCode(), result.getCode());
+
+ // db insertion fail
+ Mockito.when(workerGroupMapper.queryAllWorkerGroup()).thenReturn(Collections.singletonList(workerGroup));
+ Mockito.when(projectWorkerGroupMapper.insert(Mockito.any())).thenReturn(-1);
+ AssertionsHelper.assertThrowsServiceException(Status.ASSIGN_WORKER_GROUP_TO_PROJECT_ERROR,
+ () -> projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ getWorkerGroups()));
+
+ // success
Mockito.when(projectWorkerGroupMapper.insert(Mockito.any())).thenReturn(1);
result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
getWorkerGroups());
Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode());
+
+ // success when there is diff between current wg and assigned wg
+ Mockito.when(projectWorkerGroupMapper.selectList(Mockito.any()))
+ .thenReturn(Collections.singletonList(getDiffProjectWorkerGroup()));
+ Mockito.when(projectWorkerGroupMapper.delete(Mockito.any())).thenReturn(1);
+ result = projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ getWorkerGroups());
+ Assertions.assertEquals(Status.SUCCESS.getCode(), result.getCode());
+
+ // db deletion fail
+ Mockito.when(projectWorkerGroupMapper.delete(Mockito.any())).thenReturn(-1);
+ AssertionsHelper.assertThrowsServiceException(Status.ASSIGN_WORKER_GROUP_TO_PROJECT_ERROR,
+ () -> projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ getWorkerGroups()));
+
+ // fail when wg is referenced by task definition
+ Mockito.when(taskDefinitionMapper.queryAllDefinitionList(Mockito.anyLong()))
+ .thenReturn(Collections.singletonList(getTaskDefinitionWithDiffWorkerGroup()));
+ AssertionsHelper.assertThrowsServiceException(Status.USED_WORKER_GROUP_EXISTS,
+ () -> projectWorkerGroupRelationService.assignWorkerGroupsToProject(loginUser, projectCode,
+ getWorkerGroups()));
}
@Test
public void testQueryWorkerGroupsByProject() {
+ // no permission
+ Mockito.when(projectService.hasProjectAndPerm(Mockito.any(), Mockito.any(), Mockito.anyMap(), Mockito.any()))
+ .thenReturn(false);
+ Map result =
+ projectWorkerGroupRelationService.queryWorkerGroupsByProject(getGeneralUser(), projectCode);
+
+ Assertions.assertTrue(result.isEmpty());
+
+ // success
Mockito.when(projectService.hasProjectAndPerm(Mockito.any(), Mockito.any(), Mockito.anyMap(), Mockito.any()))
.thenReturn(true);
@@ -113,8 +178,7 @@ public class ProjectWorkerGroupRelationServiceTest {
Mockito.when(scheduleMapper.querySchedulerListByProjectName(Mockito.any()))
.thenReturn(Lists.newArrayList());
- Map result =
- projectWorkerGroupRelationService.queryWorkerGroupsByProject(getGeneralUser(), projectCode);
+ result = projectWorkerGroupRelationService.queryWorkerGroupsByProject(getGeneralUser(), projectCode);
ProjectWorkerGroup[] actualValue =
((List) result.get(Constants.DATA_LIST)).toArray(new ProjectWorkerGroup[0]);
@@ -126,20 +190,8 @@ public class ProjectWorkerGroupRelationServiceTest {
return Lists.newArrayList("default");
}
- private User getGeneralUser() {
- User loginUser = new User();
- loginUser.setUserType(UserType.GENERAL_USER);
- loginUser.setUserName("userName");
- loginUser.setId(1);
- return loginUser;
- }
-
- private User getAdminUser() {
- User loginUser = new User();
- loginUser.setUserType(UserType.ADMIN_USER);
- loginUser.setUserName("userName");
- loginUser.setId(1);
- return loginUser;
+ private List getDiffWorkerGroups() {
+ return Lists.newArrayList("default", "new");
}
private Project getProject() {
@@ -158,4 +210,20 @@ public class ProjectWorkerGroupRelationServiceTest {
projectWorkerGroup.setWorkerGroup("default");
return projectWorkerGroup;
}
+
+ private ProjectWorkerGroup getDiffProjectWorkerGroup() {
+ ProjectWorkerGroup projectWorkerGroup = new ProjectWorkerGroup();
+ projectWorkerGroup.setId(2);
+ projectWorkerGroup.setProjectCode(projectCode);
+ projectWorkerGroup.setWorkerGroup("new");
+ return projectWorkerGroup;
+ }
+
+ private TaskDefinition getTaskDefinitionWithDiffWorkerGroup() {
+ TaskDefinition taskDefinition = new TaskDefinition();
+ taskDefinition.setProjectCode(projectCode);
+ taskDefinition.setId(1);
+ taskDefinition.setWorkerGroup("new");
+ return taskDefinition;
+ }
}
diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegate.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegate.java
index 742e745fe4..620f74ef95 100644
--- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegate.java
+++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegate.java
@@ -31,12 +31,18 @@ import lombok.extern.slf4j.Slf4j;
* This class will get the property by the priority of the following: env > jvm > properties.
*/
@Slf4j
-public class ImmutablePriorityPropertyDelegate extends ImmutablePropertyDelegate {
+public class ImmutablePriorityPropertyDelegate implements IPropertyDelegate {
private static final Map>> configValueMap = new ConcurrentHashMap<>();
- public ImmutablePriorityPropertyDelegate(String propertyAbsolutePath) {
- super(propertyAbsolutePath);
+ private ImmutablePropertyDelegate immutablePropertyDelegate;
+
+ private ImmutableYamlDelegate immutableYamlDelegate;
+
+ public ImmutablePriorityPropertyDelegate(ImmutablePropertyDelegate immutablePropertyDelegate,
+ ImmutableYamlDelegate immutableYamlDelegate) {
+ this.immutablePropertyDelegate = immutablePropertyDelegate;
+ this.immutableYamlDelegate = immutableYamlDelegate;
}
@Override
@@ -56,8 +62,14 @@ public class ImmutablePriorityPropertyDelegate extends ImmutablePropertyDelegate
return value;
}
value = getConfigValueFromProperties(key);
+ if (value.isPresent()) {
+ log.debug("Get config value from properties, key: {} actualKey: {}, value: {}",
+ k, value.get().getActualKey(), value.get().getValue());
+ return value;
+ }
+ value = getConfigValueFromYaml(key);
value.ifPresent(
- stringConfigValue -> log.debug("Get config value from properties, key: {} actualKey: {}, value: {}",
+ stringConfigValue -> log.debug("Get config value from yaml, key: {} actualKey: {}, value: {}",
k, stringConfigValue.getActualKey(), stringConfigValue.getValue()));
return value;
});
@@ -76,7 +88,8 @@ public class ImmutablePriorityPropertyDelegate extends ImmutablePropertyDelegate
@Override
public Set getPropertyKeys() {
Set propertyKeys = new HashSet<>();
- propertyKeys.addAll(super.getPropertyKeys());
+ propertyKeys.addAll(this.immutablePropertyDelegate.getPropertyKeys());
+ propertyKeys.addAll(this.immutableYamlDelegate.getPropertyKeys());
propertyKeys.addAll(System.getProperties().stringPropertyNames());
propertyKeys.addAll(System.getenv().keySet());
return propertyKeys;
@@ -104,7 +117,15 @@ public class ImmutablePriorityPropertyDelegate extends ImmutablePropertyDelegate
}
private Optional> getConfigValueFromProperties(String key) {
- String value = super.get(key);
+ String value = this.immutablePropertyDelegate.get(key);
+ if (value != null) {
+ return Optional.of(ConfigValue.fromProperties(key, value));
+ }
+ return Optional.empty();
+ }
+
+ private Optional> getConfigValueFromYaml(String key) {
+ String value = this.immutableYamlDelegate.get(key);
if (value != null) {
return Optional.of(ConfigValue.fromProperties(key, value));
}
diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePropertyDelegate.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePropertyDelegate.java
index b58735afb0..4a0c192210 100644
--- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePropertyDelegate.java
+++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutablePropertyDelegate.java
@@ -49,7 +49,7 @@ public class ImmutablePropertyDelegate implements IPropertyDelegate {
} catch (IOException e) {
log.error("Load property: {} error, please check if the file exist under classpath",
propertyAbsolutePath, e);
- System.exit(1);
+ throw new RuntimeException(e);
}
}
printProperties();
diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutableYamlDelegate.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutableYamlDelegate.java
new file mode 100644
index 0000000000..5806a20fd7
--- /dev/null
+++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/config/ImmutableYamlDelegate.java
@@ -0,0 +1,82 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.common.config;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.util.Properties;
+import java.util.Set;
+
+import lombok.extern.slf4j.Slf4j;
+
+import org.springframework.beans.factory.config.YamlPropertiesFactoryBean;
+import org.springframework.core.io.InputStreamResource;
+
+@Slf4j
+public class ImmutableYamlDelegate implements IPropertyDelegate {
+
+ private static final String REMOTE_LOGGING_YAML_NAME = "/remote-logging.yaml";
+
+ private final Properties properties;
+
+ public ImmutableYamlDelegate() {
+ this(REMOTE_LOGGING_YAML_NAME);
+ }
+
+ public ImmutableYamlDelegate(String... yamlAbsolutePath) {
+ properties = new Properties();
+ // read from classpath
+ for (String fileName : yamlAbsolutePath) {
+ try (InputStream fis = ImmutableYamlDelegate.class.getResourceAsStream(fileName)) {
+ YamlPropertiesFactoryBean factory = new YamlPropertiesFactoryBean();
+ factory.setResources(new InputStreamResource(fis));
+ factory.afterPropertiesSet();
+ Properties subProperties = factory.getObject();
+ properties.putAll(subProperties);
+ } catch (IOException e) {
+ log.error("Load property: {} error, please check if the file exist under classpath",
+ yamlAbsolutePath, e);
+ throw new RuntimeException(e);
+ }
+ }
+ printProperties();
+ }
+
+ public ImmutableYamlDelegate(Properties properties) {
+ this.properties = properties;
+ }
+
+ @Override
+ public String get(String key) {
+ return properties.getProperty(key);
+ }
+
+ @Override
+ public String get(String key, String defaultValue) {
+ return properties.getProperty(key, defaultValue);
+ }
+
+ @Override
+ public Set getPropertyKeys() {
+ return properties.stringPropertyNames();
+ }
+
+ private void printProperties() {
+ properties.forEach((k, v) -> log.debug("Get property {} -> {}", k, v));
+ }
+}
diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
index ebf668a312..cc07accc9b 100644
--- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
+++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
@@ -35,6 +35,8 @@ public final class Constants {
*/
public static final String COMMON_PROPERTIES_PATH = "/common.properties";
+ public static final String REMOTE_LOGGING_YAML_PATH = "/remote-logging.yaml";
+
public static final String FORMAT_SS = "%s%s";
public static final String FORMAT_S_S = "%s/%s";
public static final String FORMAT_S_S_COLON = "%s:%s";
@@ -683,9 +685,6 @@ public final class Constants {
public static final Integer QUERY_ALL_ON_WORKFLOW = 2;
public static final Integer QUERY_ALL_ON_TASK = 3;
- /**
- * remote logging
- */
public static final String REMOTE_LOGGING_ENABLE = "remote.logging.enable";
public static final String REMOTE_LOGGING_TARGET = "remote.logging.target";
diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/PropertyUtils.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/PropertyUtils.java
index 8289b14479..82d4de9599 100644
--- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/PropertyUtils.java
+++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/PropertyUtils.java
@@ -18,9 +18,11 @@
package org.apache.dolphinscheduler.common.utils;
import static org.apache.dolphinscheduler.common.constants.Constants.COMMON_PROPERTIES_PATH;
+import static org.apache.dolphinscheduler.common.constants.Constants.REMOTE_LOGGING_YAML_PATH;
-import org.apache.dolphinscheduler.common.config.IPropertyDelegate;
import org.apache.dolphinscheduler.common.config.ImmutablePriorityPropertyDelegate;
+import org.apache.dolphinscheduler.common.config.ImmutablePropertyDelegate;
+import org.apache.dolphinscheduler.common.config.ImmutableYamlDelegate;
import java.util.HashMap;
import java.util.Map;
@@ -37,8 +39,10 @@ import com.google.common.base.Strings;
public class PropertyUtils {
// todo: add another implementation for zookeeper/etcd/consul/xx
- private static final IPropertyDelegate propertyDelegate =
- new ImmutablePriorityPropertyDelegate(COMMON_PROPERTIES_PATH);
+ private final ImmutablePriorityPropertyDelegate propertyDelegate =
+ new ImmutablePriorityPropertyDelegate(
+ new ImmutablePropertyDelegate(COMMON_PROPERTIES_PATH),
+ new ImmutableYamlDelegate(REMOTE_LOGGING_YAML_PATH));
public static String getString(String key) {
return propertyDelegate.get(key.trim());
diff --git a/dolphinscheduler-common/src/main/resources/remote-logging.yaml b/dolphinscheduler-common/src/main/resources/remote-logging.yaml
new file mode 100644
index 0000000000..2cb48750a4
--- /dev/null
+++ b/dolphinscheduler-common/src/main/resources/remote-logging.yaml
@@ -0,0 +1,61 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+remote-logging:
+ # Whether to enable remote logging
+ enable: false
+ # if remote-logging.enable = true, set the target of remote logging
+ target: OSS
+ # if remote-logging.enable = true, set the log base directory
+ base.dir: logs
+ # if remote-logging.enable = true, set the number of threads to send logs to remote storage
+ thread.pool.size: 10
+ # required if you set remote-logging.target=OSS
+ oss:
+ # oss access key id, required if you set remote-logging.target=OSS
+ access.key.id:
+ # oss access key secret, required if you set remote-logging.target=OSS
+ access.key.secret:
+ # oss bucket name, required if you set remote-logging.target=OSS
+ bucket.name:
+ # oss endpoint, required if you set remote-logging.target=OSS
+ endpoint:
+ # required if you set remote-logging.target=S3
+ s3:
+ # s3 access key id, required if you set remote-logging.target=S3
+ access.key.id:
+ # s3 access key secret, required if you set remote-logging.target=S3
+ access.key.secret:
+ # s3 bucket name, required if you set remote-logging.target=S3
+ bucket.name:
+ # s3 endpoint, required if you set remote-logging.target=S3
+ endpoint:
+ # s3 region, required if you set remote-logging.target=S3
+ region:
+ google.cloud.storage:
+ # the location of the google cloud credential, required if you set remote-logging.target=GCS
+ credential: /path/to/credential
+ # gcs bucket name, required if you set remote-logging.target=GCS
+ bucket.name:
+ abs:
+ # abs account name, required if you set resource.storage.type=ABS
+ account.name:
+ # abs account key, required if you set resource.storage.type=ABS
+ account.key:
+ # abs container name, required if you set resource.storage.type=ABS
+ container.name:
+
diff --git a/dolphinscheduler-common/src/test/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegateTest.java b/dolphinscheduler-common/src/test/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegateTest.java
index efba923a5a..6333250492 100644
--- a/dolphinscheduler-common/src/test/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegateTest.java
+++ b/dolphinscheduler-common/src/test/java/org/apache/dolphinscheduler/common/config/ImmutablePriorityPropertyDelegateTest.java
@@ -19,6 +19,7 @@ package org.apache.dolphinscheduler.common.config;
import static com.github.stefanbirkner.systemlambda.SystemLambda.withEnvironmentVariable;
import static org.apache.dolphinscheduler.common.constants.Constants.COMMON_PROPERTIES_PATH;
+import static org.apache.dolphinscheduler.common.constants.Constants.REMOTE_LOGGING_YAML_PATH;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@@ -26,7 +27,9 @@ import org.junit.jupiter.api.Test;
class ImmutablePriorityPropertyDelegateTest {
private final ImmutablePriorityPropertyDelegate immutablePriorityPropertyDelegate =
- new ImmutablePriorityPropertyDelegate(COMMON_PROPERTIES_PATH);
+ new ImmutablePriorityPropertyDelegate(
+ new ImmutablePropertyDelegate(COMMON_PROPERTIES_PATH),
+ new ImmutableYamlDelegate(REMOTE_LOGGING_YAML_PATH));
@Test
void getOverrideFromEnv() throws Exception {
diff --git a/dolphinscheduler-common/src/test/resources/remote-logging.yaml b/dolphinscheduler-common/src/test/resources/remote-logging.yaml
new file mode 100644
index 0000000000..cb149a77fe
--- /dev/null
+++ b/dolphinscheduler-common/src/test/resources/remote-logging.yaml
@@ -0,0 +1,61 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+remote-logging:
+ # Whether to enable remote logging
+ enable: false
+ # if remote-logging.enable = true, set the target of remote logging
+ target: OSS
+ # if remote-logging.enable = true, set the log base directory
+ base.dir: logs
+ # if remote-logging.enable = true, set the number of threads to send logs to remote storage
+ thread.pool.size: 10
+ # required if you set remote-logging.target=OSS
+ oss:
+ # oss access key id, required if you set remote-logging.target=OSS
+ access.key.id:
+ # oss access key secret, required if you set remote-logging.target=OSS
+ access.key.secret:
+ # oss bucket name, required if you set remote-logging.target=OSS
+ bucket.name:
+ # oss endpoint, required if you set remote-logging.target=OSS
+ endpoint:
+ # required if you set remote-logging.target=S3
+ s3:
+ # s3 access key id, required if you set remote-logging.target=S3
+ access.key.id:
+ # s3 access key secret, required if you set remote-logging.target=S3
+ access.key.secret:
+ # s3 bucket name, required if you set remote-logging.target=S3
+ bucket.name:
+ # s3 endpoint, required if you set remote-logging.target=S3
+ endpoint:
+ # s3 region, required if you set remote-logging.target=S3
+ region:
+ google.cloud.storage:
+ # the location of the google cloud credential, required if you set remote-logging.target=GCS
+ credential: /path/to/credential
+ # gcs bucket name, required if you set remote-logging.target=GCS
+ bucket.name:
+ abs:
+ # abs account name, required if you set resource.storage.type=ABS
+ account.name:
+ # abs account key, required if you set resource.storage.type=ABS
+ account.key:
+ # abs container name, required if you set resource.storage.type=ABS
+ container.name:
+
diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/CommandMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/CommandMapper.java
index 9fb6643227..8c8314e7cc 100644
--- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/CommandMapper.java
+++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/CommandMapper.java
@@ -52,14 +52,10 @@ public interface CommandMapper extends BaseMapper {
*/
List queryCommandPage(@Param("limit") int limit, @Param("offset") int offset);
- /**
- * query command page by slot
- *
- * @return command list
- */
- List queryCommandPageBySlot(@Param("limit") int limit,
- @Param("masterCount") int masterCount,
- @Param("thisMasterSlot") int thisMasterSlot);
+ List queryCommandByIdSlot(@Param("currentSlotIndex") int currentSlotIndex,
+ @Param("totalSlot") int totalSlot,
+ @Param("idStep") int idStep,
+ @Param("fetchNumber") int fetchNum);
void deleteByWorkflowInstanceIds(@Param("workflowInstanceIds") List workflowInstanceIds);
}
diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java
index 2937957dbd..664b56ee47 100644
--- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java
+++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java
@@ -56,6 +56,11 @@ public abstract class BaseDao>
return mybatisMapper.selectBatchIds(ids);
}
+ @Override
+ public List queryAll() {
+ return mybatisMapper.selectList(null);
+ }
+
@Override
public List queryByCondition(ENTITY queryCondition) {
if (queryCondition == null) {
diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/CommandDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/CommandDao.java
new file mode 100644
index 0000000000..daa52b8318
--- /dev/null
+++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/CommandDao.java
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.dao.repository;
+
+import org.apache.dolphinscheduler.dao.entity.Command;
+
+import java.util.List;
+
+public interface CommandDao extends IDao {
+
+ /**
+ * Query command by command id and server slot, return the command which match (commandId / step) %s totalSlot = currentSlotIndex
+ *
+ * @param currentSlotIndex current slot index
+ * @param totalSlot total slot number
+ * @param idStep id step in db
+ * @param fetchNum fetch number
+ * @return command list
+ */
+ List queryCommandByIdSlot(int currentSlotIndex,
+ int totalSlot,
+ int idStep,
+ int fetchNum);
+}
diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java
index c566d9b904..ab77419600 100644
--- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java
+++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java
@@ -41,6 +41,11 @@ public interface IDao {
*/
List queryByIds(Collection extends Serializable> ids);
+ /**
+ * Query all entities.
+ */
+ List queryAll();
+
/**
* Query the entity by condition.
*/
diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/CommandDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/CommandDaoImpl.java
new file mode 100644
index 0000000000..0b510d15b5
--- /dev/null
+++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/CommandDaoImpl.java
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.dolphinscheduler.dao.repository.impl;
+
+import org.apache.dolphinscheduler.dao.entity.Command;
+import org.apache.dolphinscheduler.dao.mapper.CommandMapper;
+import org.apache.dolphinscheduler.dao.repository.BaseDao;
+import org.apache.dolphinscheduler.dao.repository.CommandDao;
+
+import java.util.List;
+
+import org.springframework.stereotype.Repository;
+
+@Repository
+public class CommandDaoImpl extends BaseDao implements CommandDao {
+
+ public CommandDaoImpl(CommandMapper commandMapper) {
+ super(commandMapper);
+ }
+
+ @Override
+ public List queryCommandByIdSlot(int currentSlotIndex, int totalSlot, int idStep, int fetchNum) {
+ return mybatisMapper.queryCommandByIdSlot(currentSlotIndex, totalSlot, idStep, fetchNum);
+ }
+
+}
diff --git a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/CommandMapper.xml b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/CommandMapper.xml
index 56db890ef0..16f7c05f25 100644
--- a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/CommandMapper.xml
+++ b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/CommandMapper.xml
@@ -40,12 +40,12 @@
limit #{limit} offset #{offset}
-