Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
三久
DolphinScheduler
提交
bd0fa79e
DolphinScheduler
项目概览
三久
/
DolphinScheduler
与 Fork 源项目一致
Fork自
apache / DolphinScheduler
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
DolphinScheduler
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
未验证
提交
bd0fa79e
编写于
9月 15, 2022
作者:
S
Stalary
提交者:
GitHub
9月 15, 2022
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[Bug][Dependent]: Id also clone due to duplicate when use dependent mode. (#11929)
上级
e938fdbe
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
65 addition
and
5 deletion
+65
-5
dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java
...olphinscheduler/api/service/impl/ExecutorServiceImpl.java
+4
-2
dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ExecutorServiceTest.java
...che/dolphinscheduler/api/service/ExecutorServiceTest.java
+61
-3
未找到文件。
dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java
浏览文件 @
bd0fa79e
...
...
@@ -887,7 +887,7 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ
/**
* create complement dependent command
*/
p
rotected
int
createComplementDependentCommand
(
List
<
Schedule
>
schedules
,
Command
command
)
{
p
ublic
int
createComplementDependentCommand
(
List
<
Schedule
>
schedules
,
Command
command
)
{
int
dependentProcessDefinitionCreateCount
=
0
;
Command
dependentCommand
;
...
...
@@ -901,9 +901,11 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ
List
<
DependentProcessDefinition
>
dependentProcessDefinitionList
=
getComplementDependentDefinitionList
(
dependentCommand
.
getProcessDefinitionCode
(),
CronUtils
.
getMaxCycle
(
schedules
.
get
(
0
).
getCrontab
()),
dependentCommand
.
getWorkerGroup
());
dependentCommand
.
setTaskDependType
(
TaskDependType
.
TASK_POST
);
for
(
DependentProcessDefinition
dependentProcessDefinition
:
dependentProcessDefinitionList
)
{
// If the id is Integer, the auto-increment id will be obtained by mybatis-plus
// and causing duplicate when clone it.
dependentCommand
.
setId
(
null
);
dependentCommand
.
setProcessDefinitionCode
(
dependentProcessDefinition
.
getProcessDefinitionCode
());
dependentCommand
.
setProcessDefinitionVersion
(
dependentProcessDefinition
.
getProcessDefinitionVersion
());
dependentCommand
.
setWorkerGroup
(
dependentProcessDefinition
.
getWorkerGroup
());
...
...
dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ExecutorServiceTest.java
浏览文件 @
bd0fa79e
...
...
@@ -20,18 +20,30 @@ package org.apache.dolphinscheduler.api.service;
import
static
org
.
apache
.
dolphinscheduler
.
api
.
constants
.
ApiFuncIdentificationConstant
.
RERUN
;
import
static
org
.
apache
.
dolphinscheduler
.
api
.
constants
.
ApiFuncIdentificationConstant
.
WORKFLOW_START
;
import
static
org
.
mockito
.
ArgumentMatchers
.
any
;
import
static
org
.
mockito
.
ArgumentMatchers
.
argThat
;
import
static
org
.
mockito
.
Mockito
.
doReturn
;
import
static
org
.
mockito
.
Mockito
.
times
;
import
static
org
.
mockito
.
Mockito
.
verify
;
import
org.apache.dolphinscheduler.api.enums.ExecuteType
;
import
org.apache.dolphinscheduler.api.enums.Status
;
import
org.apache.dolphinscheduler.api.permission.ResourcePermissionCheckService
;
import
org.apache.dolphinscheduler.api.service.impl.BaseServiceImpl
;
import
org.apache.dolphinscheduler.api.service.impl.ExecutorServiceImpl
;
import
org.apache.dolphinscheduler.api.service.impl.ProjectServiceImpl
;
import
org.apache.dolphinscheduler.common.Constants
;
import
org.apache.dolphinscheduler.common.enums.*
;
import
org.apache.dolphinscheduler.common.enums.CommandType
;
import
org.apache.dolphinscheduler.common.enums.ComplementDependentMode
;
import
org.apache.dolphinscheduler.common.enums.FailureStrategy
;
import
org.apache.dolphinscheduler.common.enums.Priority
;
import
org.apache.dolphinscheduler.common.enums.ReleaseState
;
import
org.apache.dolphinscheduler.common.enums.RunMode
;
import
org.apache.dolphinscheduler.common.enums.TaskGroupQueueStatus
;
import
org.apache.dolphinscheduler.common.enums.WarningType
;
import
org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus
;
import
org.apache.dolphinscheduler.common.model.Server
;
import
org.apache.dolphinscheduler.dao.entity.Command
;
import
org.apache.dolphinscheduler.dao.entity.DependentProcessDefinition
;
import
org.apache.dolphinscheduler.dao.entity.ProcessDefinition
;
import
org.apache.dolphinscheduler.dao.entity.ProcessInstance
;
import
org.apache.dolphinscheduler.dao.entity.Project
;
...
...
@@ -45,17 +57,18 @@ import org.apache.dolphinscheduler.dao.mapper.ProcessTaskRelationMapper;
import
org.apache.dolphinscheduler.dao.mapper.ProjectMapper
;
import
org.apache.dolphinscheduler.dao.mapper.TaskDefinitionMapper
;
import
org.apache.dolphinscheduler.dao.mapper.TaskGroupQueueMapper
;
import
org.apache.dolphinscheduler.api.permission.ResourcePermissionCheckService
;
import
org.apache.dolphinscheduler.service.process.ProcessService
;
import
java.util.ArrayList
;
import
java.util.Collections
;
import
java.util.Date
;
import
java.util.HashMap
;
import
java.util.LinkedList
;
import
java.util.List
;
import
java.util.Map
;
import
java.util.Optional
;
import
org.assertj.core.util.Lists
;
import
org.junit.Assert
;
import
org.junit.Before
;
import
org.junit.Test
;
...
...
@@ -178,7 +191,8 @@ public class ExecutorServiceTest {
.
thenReturn
(
checkProjectAndAuth
());
Mockito
.
when
(
processDefinitionMapper
.
queryByCode
(
processDefinitionCode
)).
thenReturn
(
processDefinition
);
Mockito
.
when
(
processService
.
getTenantForProcess
(
tenantId
,
userId
)).
thenReturn
(
new
Tenant
());
Mockito
.
when
(
processService
.
createCommand
(
any
(
Command
.
class
))).
thenReturn
(
1
);
doReturn
(
1
).
when
(
processService
).
createCommand
(
argThat
(
c
->
c
.
getId
()
==
null
));
doReturn
(
0
).
when
(
processService
).
createCommand
(
argThat
(
c
->
c
.
getId
()
!=
null
));
Mockito
.
when
(
monitorService
.
getServerListFromRegistry
(
true
)).
thenReturn
(
getMasterServersList
());
Mockito
.
when
(
processService
.
findProcessInstanceDetailById
(
processInstanceId
)).
thenReturn
(
Optional
.
ofNullable
(
processInstance
));
Mockito
.
when
(
processService
.
findProcessDefinition
(
1L
,
1
)).
thenReturn
(
processDefinition
);
...
...
@@ -237,6 +251,50 @@ public class ExecutorServiceTest {
}
@Test
public
void
testComplementWithDependentMode
()
{
Schedule
schedule
=
new
Schedule
();
schedule
.
setStartTime
(
new
Date
());
schedule
.
setEndTime
(
new
Date
());
schedule
.
setCrontab
(
"0 0 7 * * ? *"
);
schedule
.
setFailureStrategy
(
FailureStrategy
.
CONTINUE
);
schedule
.
setReleaseState
(
ReleaseState
.
OFFLINE
);
schedule
.
setWarningType
(
WarningType
.
NONE
);
schedule
.
setCreateTime
(
new
Date
());
schedule
.
setUpdateTime
(
new
Date
());
List
<
Schedule
>
schedules
=
Lists
.
newArrayList
(
schedule
);
Mockito
.
when
(
processService
.
queryReleaseSchedulerListByProcessDefinitionCode
(
processDefinitionCode
))
.
thenReturn
(
schedules
);
DependentProcessDefinition
dependentProcessDefinition
=
new
DependentProcessDefinition
();
dependentProcessDefinition
.
setProcessDefinitionCode
(
2
);
dependentProcessDefinition
.
setProcessDefinitionVersion
(
1
);
dependentProcessDefinition
.
setTaskDefinitionCode
(
1
);
dependentProcessDefinition
.
setWorkerGroup
(
Constants
.
DEFAULT_WORKER_GROUP
);
dependentProcessDefinition
.
setTaskParams
(
"{\"localParams\":[],\"resourceList\":[],\"dependence\":{\"relation\":\"AND\",\"dependTaskList\":[{\"relation\":\"AND\",\"dependItemList\":[{\"depTaskCode\":2,\"status\":\"SUCCESS\"}]}]},\"conditionResult\":{\"successNode\":[1],\"failedNode\":[1]}}"
);
Mockito
.
when
(
processService
.
queryDependentProcessDefinitionByProcessDefinitionCode
(
processDefinitionCode
))
.
thenReturn
(
Lists
.
newArrayList
(
dependentProcessDefinition
));
Map
<
Long
,
String
>
processDefinitionWorkerGroupMap
=
new
HashMap
<>();
processDefinitionWorkerGroupMap
.
put
(
1L
,
Constants
.
DEFAULT_WORKER_GROUP
);
Mockito
.
when
(
processService
.
queryWorkerGroupByProcessDefinitionCodes
(
Lists
.
newArrayList
(
1L
)))
.
thenReturn
(
processDefinitionWorkerGroupMap
);
Command
command
=
new
Command
();
command
.
setId
(
1
);
command
.
setCommandType
(
CommandType
.
COMPLEMENT_DATA
);
command
.
setCommandParam
(
"{\"StartNodeList\":\"1\",\"complementStartDate\":\"2020-01-01 00:00:00\",\"complementEndDate\":\"2020-01-31 23:00:00\"}"
);
command
.
setWorkerGroup
(
Constants
.
DEFAULT_WORKER_GROUP
);
command
.
setProcessDefinitionCode
(
processDefinitionCode
);
command
.
setExecutorId
(
1
);
int
count
=
executorService
.
createComplementDependentCommand
(
schedules
,
command
);
Assert
.
assertEquals
(
1
,
count
);
}
/**
* date error
*/
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录