Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
B
byit-myth-job
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
liyuan
byit-myth-job
Commits
e0279415
Commit
e0279415
authored
Mar 17, 2020
by
guo_minglei@163.com
Browse files
Options
Browse Files
Download
Plain Diff
Merge remote-tracking branch 'origin/developer' into developer
parents
9c56b8a4
4d349a72
Show whitespace changes
Inline
Side-by-side
Showing
5 changed files
with
371 additions
and
11 deletions
+371
-11
logback-spring.xml
byit-myth-admin/src/main/resources/logback-spring.xml
+3
-3
RunRecordingAndJobTaskServiceImpl.java
...ce/mapservice/impl/RunRecordingAndJobTaskServiceImpl.java
+21
-1
DateUtil.java
...ore-common/src/main/java/com/byit/job/utils/DateUtil.java
+2
-2
Test1.java
...core/myth-executor-core/src/test/java/com/test/Test1.java
+78
-5
AddComplexPy.java
...-demo-client/src/main/java/com/byit/job/AddComplexPy.java
+267
-0
No files found.
byit-myth-admin/src/main/resources/logback-spring.xml
View file @
e0279415
...
@@ -58,13 +58,13 @@
...
@@ -58,13 +58,13 @@
</encoder>
</encoder>
</appender>
</appender>
<logger
name=
"com.byit.thread"
level=
"
debug
"
additivity=
"false"
>
<logger
name=
"com.byit.thread"
level=
"
info
"
additivity=
"false"
>
<appender-ref
ref=
"console"
/>
<appender-ref
ref=
"console"
/>
</logger>
</logger>
<logger
name=
"com.byit.selector"
level=
"
debug
"
additivity=
"false"
>
<logger
name=
"com.byit.selector"
level=
"
info
"
additivity=
"false"
>
<appender-ref
ref=
"console"
/>
<appender-ref
ref=
"console"
/>
</logger>
</logger>
<logger
name=
"com.byit.task"
level=
"
debug
"
additivity=
"false"
>
<logger
name=
"com.byit.task"
level=
"
info
"
additivity=
"false"
>
<appender-ref
ref=
"console"
/>
<appender-ref
ref=
"console"
/>
</logger>
</logger>
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/mapservice/impl/RunRecordingAndJobTaskServiceImpl.java
View file @
e0279415
package
com
.
byit
.
service
.
mapservice
.
impl
;
package
com
.
byit
.
service
.
mapservice
.
impl
;
import
cn.hutool.core.collection.CollectionUtil
;
import
com.byit.enums.EmailEnum
;
import
com.byit.enums.EmailEnum
;
import
com.byit.enums.FlowPropertyEnum
;
import
com.byit.enums.FlowPropertyEnum
;
import
com.byit.enums.NodeRunStatusPropertyEnum
;
import
com.byit.enums.NodeRunStatusPropertyEnum
;
import
com.byit.enums.RunRecordingEnum
;
import
com.byit.enums.RunRecordingEnum
;
import
com.byit.job.utils.DateUtil
;
import
com.byit.model.*
;
import
com.byit.model.*
;
import
com.byit.service.*
;
import
com.byit.service.*
;
import
com.byit.service.mapservice.RunRecordingAndJobTaskService
;
import
com.byit.service.mapservice.RunRecordingAndJobTaskService
;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.commons.lang3.StringUtils
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
import
org.springframework.transaction.annotation.Transactional
;
import
org.springframework.transaction.annotation.Transactional
;
...
@@ -33,10 +36,14 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -33,10 +36,14 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
private
final
JobTaskService
jobTaskService
;
private
final
JobTaskService
jobTaskService
;
private
final
WaitingRecordService
waitingRecordService
;
private
final
WaitingRecordService
waitingRecordService
;
private
final
WaitingTaskService
taskService
;
private
final
WaitingTaskService
taskService
;
/**
* 节点依赖查询操作
*/
private
final
NodeDependencyService
nodeDependencyService
;
public
RunRecordingAndJobTaskServiceImpl
(
JobTaskRunLogService
jobTaskRunLogService
,
NodeService
nodeService
,
public
RunRecordingAndJobTaskServiceImpl
(
JobTaskRunLogService
jobTaskRunLogService
,
NodeService
nodeService
,
FlowService
flowService
,
RunRecordingService
runRecordingService
,
FlowService
flowService
,
RunRecordingService
runRecordingService
,
JobTaskService
jobTaskService
,
WaitingRecordService
waitingRecordService
,
WaitingTaskService
taskService
)
{
JobTaskService
jobTaskService
,
WaitingRecordService
waitingRecordService
,
WaitingTaskService
taskService
,
NodeDependencyService
nodeDependencyService
)
{
this
.
jobTaskRunLogService
=
jobTaskRunLogService
;
this
.
jobTaskRunLogService
=
jobTaskRunLogService
;
this
.
nodeService
=
nodeService
;
this
.
nodeService
=
nodeService
;
this
.
flowService
=
flowService
;
this
.
flowService
=
flowService
;
...
@@ -44,6 +51,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -44,6 +51,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
this
.
jobTaskService
=
jobTaskService
;
this
.
jobTaskService
=
jobTaskService
;
this
.
waitingRecordService
=
waitingRecordService
;
this
.
waitingRecordService
=
waitingRecordService
;
this
.
taskService
=
taskService
;
this
.
taskService
=
taskService
;
this
.
nodeDependencyService
=
nodeDependencyService
;
}
}
/**
/**
* 保存节点日志
* 保存节点日志
...
@@ -100,6 +108,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -100,6 +108,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
.
triggerTime
(
equals
?
mainFlow
.
getTriggerNextTime
():
virFlow
.
getTriggerNextTime
())
.
triggerTime
(
equals
?
mainFlow
.
getTriggerNextTime
():
virFlow
.
getTriggerNextTime
())
.
principal
(
virFlow
.
getPrincipal
())
.
principal
(
virFlow
.
getPrincipal
())
.
startTime
(
new
Date
())
.
startTime
(
new
Date
())
.
flowNodeCount
(
virFlow
.
getFlowNodeCount
())
.
isAlarm
(
EmailEnum
.
IS_ALARM_NO
.
getCode
())
.
isAlarm
(
EmailEnum
.
IS_ALARM_NO
.
getCode
())
.
isInner
(
FlowPropertyEnum
.
IS_INNER
.
getCode
())
.
isInner
(
FlowPropertyEnum
.
IS_INNER
.
getCode
())
.
failFast
(
RunRecordingEnum
.
FAIL_FAST_NO
.
getCode
()).
build
();
.
failFast
(
RunRecordingEnum
.
FAIL_FAST_NO
.
getCode
()).
build
();
...
@@ -110,6 +119,12 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -110,6 +119,12 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
List
<
JobTask
>
jobTasks
=
nodeByFlowIdAndVersionName
.
stream
()
List
<
JobTask
>
jobTasks
=
nodeByFlowIdAndVersionName
.
stream
()
.
map
(
node
->
{
.
map
(
node
->
{
JobTask
task
=
new
JobTask
();
JobTask
task
=
new
JobTask
();
List
<
Integer
>
dependIdByNodeId
=
nodeDependencyService
.
findDependIdByNodeId
(
jobTask
.
getNodeId
());
if
(
CollectionUtil
.
isNotEmpty
(
dependIdByNodeId
)){
String
parentIds
=
StringUtils
.
join
(
dependIdByNodeId
,
","
);
task
.
setNodeDepend
(
parentIds
);
}
BeanUtils
.
copyProperties
(
node
,
task
);
BeanUtils
.
copyProperties
(
node
,
task
);
if
(
equals
)
{
if
(
equals
)
{
task
.
setTriggerTime
(
mainFlow
.
getTriggerNextTime
());
task
.
setTriggerTime
(
mainFlow
.
getTriggerNextTime
());
...
@@ -117,6 +132,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -117,6 +132,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
task
.
setTriggerTime
(
virFlow
.
getTriggerNextTime
());
task
.
setTriggerTime
(
virFlow
.
getTriggerNextTime
());
}
}
task
.
setRunId
(
jobTask
.
getRunId
());
task
.
setRunId
(
jobTask
.
getRunId
());
task
.
setTriggerStatus
(
"1"
);
return
task
;
return
task
;
}).
collect
(
Collectors
.
toList
());
}).
collect
(
Collectors
.
toList
());
jobTaskService
.
saveJobTasks
(
jobTasks
);
jobTaskService
.
saveJobTasks
(
jobTasks
);
...
@@ -148,4 +164,8 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -148,4 +164,8 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
waitingRecordService
.
updateById
(
waitingRecord
);
waitingRecordService
.
updateById
(
waitingRecord
);
}
}
public
static
void
main
(
String
[]
args
)
{
System
.
out
.
println
(
DateUtil
.
dateFormat
(
new
Date
(
1584420360000L
),
"yyyy-MM-dd HH mm ss"
));
}
}
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/utils/DateUtil.java
View file @
e0279415
...
@@ -25,9 +25,9 @@ public class DateUtil {
...
@@ -25,9 +25,9 @@ public class DateUtil {
LocalDateTime
localDateTime
=
dateToLocalDateTime
(
date
);
LocalDateTime
localDateTime
=
dateToLocalDateTime
(
date
);
String
dateFormat
;
String
dateFormat
;
if
(
StringUtils
.
isNotBlank
(
formatStr
)){
if
(
StringUtils
.
isNotBlank
(
formatStr
)){
dateFormat
=
localDateTime
.
format
(
DateTimeFormatter
.
ofPattern
(
NOT_FORMAT_DATE
));
}
else
{
dateFormat
=
localDateTime
.
format
(
DateTimeFormatter
.
ofPattern
(
formatStr
));
dateFormat
=
localDateTime
.
format
(
DateTimeFormatter
.
ofPattern
(
formatStr
));
}
else
{
dateFormat
=
localDateTime
.
format
(
DateTimeFormatter
.
ofPattern
(
NOT_FORMAT_DATE
));
}
}
return
dateFormat
;
return
dateFormat
;
}
}
...
...
byit-myth-core/myth-executor-core/src/test/java/com/test/Test1.java
View file @
e0279415
...
@@ -13,10 +13,83 @@ import java.util.Map;
...
@@ -13,10 +13,83 @@ import java.util.Map;
public
class
Test1
{
public
class
Test1
{
public
static
void
main
(
String
[]
args
)
throws
IOException
,
MyException
{
public
static
void
main
(
String
[]
args
)
throws
IOException
,
MyException
{
FastDfsFileSystem
fds
=
new
FastDfsFileSystem
();
FastDfsFileSystem
fds
=
new
FastDfsFileSystem
();
fds
.
fileRemove
(
"ddmp/M00/00/00/CgB4Al5XoBSAZg1qAAAAW4eTPeg3226.py"
);
/*fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyeAZdkbAAAAb-omIFM6333.py");
byte
[]
bytes
=
FileUtils
.
readFileToByteArray
(
new
File
(
"D:\\2020project\\byit-myth-job\\byit-myth-core\\myth-executor-core\\src\\test\\java\\com\\test\\test-ex.py"
));
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyeANch6AAAAb1xkZcY9815.py");
Map
<
String
,
String
>
map
=
new
HashMap
<>();
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyeAY9HjAAAAbzGlprU0241.py");
map
.
put
(
"filename"
,
"test-ex.py"
);
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyeASM74AAAAb-uR6K03377.py");
System
.
out
.
println
(
fds
.
uploadFile
(
bytes
,
"py"
,
map
));
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiAPlYXAAAAb4ZQK943410.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiAc_nJAAAAbzASbks9878.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiABB-kAAAAb13TrTg1071.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiARAFVAAAAb18L9Do3337.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiACVSEAAAAbzLKN0k8232.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiAZbfWAAAAcK8DUJU5054.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wNyiAZnPBAAAAcMLCk-Y6177.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wODWAePJQAAAAHSB_HGI4727.py");
fds.fileRemove("ddmp/M00/00/00/CgB4Al5wODWAFdLWAAAAG7nO0c08696.py");*/
/*byte[] bytes1 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node1.py"));
Map<String,String> map1 = new HashMap<>();
map1.put("filename","node1.py");
System.out.println(fds.uploadFile(bytes1, "py", map1));
byte[] bytes2 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node2.py"));
Map<String,String> map2 = new HashMap<>();
map2.put("filename","node2.py");
System.out.println(fds.uploadFile(bytes2, "py", map2));
byte[] bytes3 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node3.py"));
Map<String,String> map3 = new HashMap<>();
map3.put("filename","node3.py");
System.out.println(fds.uploadFile(bytes3, "py", map3));
byte[] bytes4 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node4.py"));
Map<String,String> map4 = new HashMap<>();
map4.put("filename","node4.py");
System.out.println(fds.uploadFile(bytes4, "py", map4));
byte[] bytes5 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node5.py"));
Map<String,String> map5 = new HashMap<>();
map5.put("filename","node5.py");
System.out.println(fds.uploadFile(bytes5, "py", map5));
byte[] bytes6 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node6.py"));
Map<String,String> map6 = new HashMap<>();
map6.put("filename","node6.py");
System.out.println(fds.uploadFile(bytes6, "py", map6));
byte[] bytes7 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node7.py"));
Map<String,String> map7 = new HashMap<>();
map7.put("filename","node7.py");
System.out.println(fds.uploadFile(bytes7, "py", map7));
byte[] bytes8 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node8.py"));
Map<String,String> map8 = new HashMap<>();
map8.put("filename","node8.py");
System.out.println(fds.uploadFile(bytes8, "py", map8));
byte[] bytes9 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node9.py"));
Map<String,String> map9 = new HashMap<>();
map9.put("filename","node9.py");
System.out.println(fds.uploadFile(bytes9, "py", map9));
byte[] bytes10 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node10.py"));
Map<String,String> map10 = new HashMap<>();
map10.put("filename","node10.py");
System.out.println(fds.uploadFile(bytes10, "py", map10));
byte[] bytes11 = FileUtils.readFileToByteArray(new File("C:\\Users\\Administrator\\Desktop\\script/node11.py"));
Map<String,String> map11 = new HashMap<>();
map11.put("filename","node11.py");
System.out.println(fds.uploadFile(bytes11, "py", map11));*/
byte
[]
bytes1
=
FileUtils
.
readFileToByteArray
(
new
File
(
"C:\\Users\\Administrator\\Desktop\\script/start.py"
));
Map
<
String
,
String
>
map1
=
new
HashMap
<>();
map1
.
put
(
"filename"
,
"start.py"
);
System
.
out
.
println
(
fds
.
uploadFile
(
bytes1
,
"py"
,
map1
));
byte
[]
bytes2
=
FileUtils
.
readFileToByteArray
(
new
File
(
"C:\\Users\\Administrator\\Desktop\\script/end.py"
));
Map
<
String
,
String
>
map2
=
new
HashMap
<>();
map2
.
put
(
"filename"
,
"end.py"
);
System
.
out
.
println
(
fds
.
uploadFile
(
bytes2
,
"py"
,
map2
));
}
}
}
}
demo-client/byit-demo-client/src/main/java/com/byit/job/AddComplexPy.java
0 → 100644
View file @
e0279415
package
com
.
byit
.
job
;
import
com.byit.job.dto.ScriptParamAndPlaceholderDto
;
import
com.byit.job.dto.plugin.*
;
import
com.byit.utils.JobUtils
;
import
java.util.*
;
import
java.util.concurrent.TimeUnit
;
/**
* @author huanfgu
*/
public
class
AddComplexPy
{
public
static
void
main
(
String
[]
args
)
{
PluginPackage
pluginPackage
=
new
PluginPackage
();
pluginPackage
.
setWorkspaceName
(
"test"
);
pluginPackage
.
setFlow
(
createFlow
());
JobUtils
.
publish
(
pluginPackage
);
}
public
static
PluginFlow
createFlow
(){
PluginFlow
pluginFlow
=
new
PluginFlow
();
//构建工作流信息
PluginFlowConfig
build
=
PluginFlowConfig
.
builder
().
alarmEmail
(
"huangfukexing@byitgroup.com"
)
.
alarmlAction
(
"1"
)
.
execType
(
"1"
)
.
flowCron
(
"0 0/1 * * * ? *"
)
.
flowTimeout
(
TimeUnit
.
MINUTES
.
toMillis
(
30
))
.
priority
(
"2"
)
.
repeatCount
(
1
)
.
scheduleFollow
(
"1"
)
.
build
();
pluginFlow
.
setName
(
"复杂工作流"
);
pluginFlow
.
setDesc
(
"测试多脚本复杂工作流创建"
);
pluginFlow
.
setConfig
(
build
);
pluginFlow
.
setPrincipal
(
"皇甫科星"
);
pluginFlow
.
setRePublish
(
false
);
pluginFlow
.
setAuthor
(
"huangfukexing"
);
pluginFlow
.
setNodeList
(
createNodes
());
return
pluginFlow
;
}
/**
* 创建节点
* @return
*/
public
static
List
<
PluginBaseNode
>
createNodes
(){
//构建开始节点
PluginNode
startNode
=
new
PluginNode
();
PluginNodeConfig
startConf
=
new
PluginNodeConfig
();
nodeSet
(
startNode
,
false
);
startNode
.
setName
(
"start"
);
startNode
.
setDesc
(
"开始节点"
);
startNode
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa6WAfhQ9AAAAN9Xm1mU1221.py"
);
confSet
(
startConf
);
startNode
.
setConfig
(
startConf
);
System
.
out
.
println
(
"-----------------start 节点构建成功,开始构建node1节点,依赖start节点--------------------"
);
PluginNode
node1
=
new
PluginNode
();
PluginNodeConfig
node1Conf
=
new
PluginNodeConfig
();
nodeSet
(
node1
,
true
);
node1
.
setName
(
"node1"
);
node1
.
setDesc
(
"node1节点"
);
node1
.
setScriptUrls
(
"ddmp/M00/00/00/CgB4Al5wa36AaUJ7AAAAiYy1k-k9801.py"
);
confSet
(
node1Conf
);
node1
.
setConfig
(
node1Conf
);
node1
.
setDependNodeNameList
(
Collections
.
singletonList
(
"start"
));
System
.
out
.
println
(
"-----------------node1 节点构建成功,开始构建 node2节点,依赖start节点--------------------"
);
PluginNode
node2
=
new
PluginNode
();
PluginNodeConfig
node2Conf
=
new
PluginNodeConfig
();
nodeSet
(
node2
,
true
);
node2
.
setName
(
"node2"
);
node2
.
setDesc
(
"node2节点"
);
node2
.
setScriptUrls
(
"ddmp/M00/00/00/CgB4Al5wa36AScI_AAAAiTr31nw2468.py"
);
confSet
(
node2Conf
);
node2
.
setConfig
(
node2Conf
);
node2
.
setDependNodeNameList
(
Collections
.
singletonList
(
"start"
));
System
.
out
.
println
(
"-----------------node2 节点构建成功,开始构建 node3节点,依赖node1 node2节点--------------------"
);
PluginNode
node3
=
new
PluginNode
();
PluginNodeConfig
node3Conf
=
new
PluginNodeConfig
();
nodeSet
(
node3
,
true
);
node3
.
setName
(
"node3"
);
node3
.
setDesc
(
"node3节点"
);
node3
.
setScriptUrls
(
"ddmp/M00/00/00/CgB4Al5wa36ADcUtAAAAiVc2FQ85978.py"
);
confSet
(
node3Conf
);
node3
.
setConfig
(
node3Conf
);
node3
.
setDependNodeNameList
(
Arrays
.
asList
(
"node1"
,
"node2"
));
System
.
out
.
println
(
"-----------------node3 节点构建成功,开始构建 node4 节点(虚节点),依赖node3节点--------------------"
);
PluginFlow
node4In
=
createInFlow
();
node4In
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node3"
));
System
.
out
.
println
(
"-----------node4 节点构建成功,开始构建 node10 节点,依赖node4节点(虚节点)----------------"
);
PluginNode
node10
=
new
PluginNode
();
PluginNodeConfig
node10Conf
=
new
PluginNodeConfig
();
nodeSet
(
node10
,
true
);
node10
.
setName
(
"node10"
);
node10
.
setDesc
(
"node10节点"
);
node10
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36ALFnFAAAAioTRmbQ1384.py"
);
confSet
(
node10Conf
);
node10
.
setConfig
(
node10Conf
);
node10
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node4"
));
System
.
out
.
println
(
"-----------node10 节点构建成功,开始构建 node11 节点,依赖node4节点(虚节点)----------------"
);
PluginNode
node11
=
new
PluginNode
();
PluginNodeConfig
node11Conf
=
new
PluginNodeConfig
();
nodeSet
(
node11
,
true
);
node11
.
setName
(
"node11"
);
node11
.
setDesc
(
"node11节点"
);
node11
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36AOM7_AAAAiukQWsc6165.py"
);
confSet
(
node11Conf
);
node11
.
setConfig
(
node11Conf
);
node11
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node4"
));
System
.
out
.
println
(
"-----------node11 节点构建成功,开始构建 end 节点,依赖 node10 node11 节点----------------"
);
PluginNode
endNode
=
new
PluginNode
();
PluginNodeConfig
endNodeConf
=
new
PluginNodeConfig
();
nodeSet
(
endNode
,
false
);
endNode
.
setName
(
"end"
);
endNode
.
setDesc
(
"end节点"
);
endNode
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa6WADyyrAAAANb42CIc7700.py"
);
confSet
(
endNodeConf
);
endNode
.
setConfig
(
endNodeConf
);
endNode
.
setDependNodeNameList
(
Arrays
.
asList
(
"node10"
,
"node11"
));
return
Arrays
.
asList
(
startNode
,
node1
,
node2
,
node3
,
node4In
,
node10
,
node11
,
endNode
);
}
/**
* 创建虚节点工作流
* @return
*/
public
static
PluginFlow
createInFlow
(){
PluginFlow
pluginFlow
=
new
PluginFlow
();
PluginNode
startNode
=
new
PluginNode
();
PluginNodeConfig
startConf
=
new
PluginNodeConfig
();
nodeSet
(
startNode
,
false
);
startNode
.
setName
(
"start"
);
startNode
.
setDesc
(
"开始节点"
);
startNode
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa6WAfhQ9AAAAN9Xm1mU1221.py"
);
confSet
(
startConf
);
startNode
.
setConfig
(
startConf
);
System
.
out
.
println
(
"-----------------start 节点构建成功,开始构建node4节点,依赖start节点--------------------"
);
PluginNode
node4
=
new
PluginNode
();
PluginNodeConfig
node4Conf
=
new
PluginNodeConfig
();
nodeSet
(
node4
,
true
);
node4
.
setName
(
"node4"
);
node4
.
setDesc
(
"node4节点"
);
node4
.
setScriptUrls
(
"ddmp/M00/00/00/CgB4Al5wa36ARJKOAAAAiY0CWxc8942.py"
);
confSet
(
node4Conf
);
node4
.
setConfig
(
node4Conf
);
node4
.
setDependNodeNameList
(
Collections
.
singletonList
(
"start"
));
System
.
out
.
println
(
"-----------------node4 节点构建成功,开始构建node5节点,依赖 node4 节点--------------------"
);
PluginNode
node5
=
new
PluginNode
();
PluginNodeConfig
node5Conf
=
new
PluginNodeConfig
();
nodeSet
(
node5
,
true
);
node5
.
setName
(
"node5"
);
node5
.
setDesc
(
"node5节点"
);
node5
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36AKg3oAAAAieDDmGQ1351.py"
);
confSet
(
node5Conf
);
node5
.
setConfig
(
node5Conf
);
node5
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node4"
));
System
.
out
.
println
(
"-----------------node5 节点构建成功,开始构建node6节点,依赖 node4 节点--------------------"
);
PluginNode
node6
=
new
PluginNode
();
PluginNodeConfig
node6Conf
=
new
PluginNodeConfig
();
nodeSet
(
node6
,
true
);
node6
.
setName
(
"node6"
);
node6
.
setDesc
(
"node6节点"
);
node6
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36ATQGXAAAAiVaB3fE3326.py"
);
confSet
(
node6Conf
);
node6
.
setConfig
(
node6Conf
);
node6
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node4"
));
System
.
out
.
println
(
"-----------------node6 节点构建成功,开始构建node7节点,依赖 node5 节点--------------------"
);
PluginNode
node7
=
new
PluginNode
();
PluginNodeConfig
node7Conf
=
new
PluginNodeConfig
();
nodeSet
(
node7
,
true
);
node7
.
setName
(
"node7"
);
node7
.
setDesc
(
"node7节点"
);
node7
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36AFC91AAAAiTtAHoI8044.py"
);
confSet
(
node7Conf
);
node7
.
setConfig
(
node7Conf
);
node7
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node5"
));
System
.
out
.
println
(
"-----------------node7 节点构建成功,开始构建node8节点,依赖 node6 节点--------------------"
);
PluginNode
node8
=
new
PluginNode
();
PluginNodeConfig
node8Conf
=
new
PluginNodeConfig
();
nodeSet
(
node8
,
true
);
node8
.
setName
(
"node8"
);
node8
.
setDesc
(
"node8节点"
);
node8
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36AeBhfAAAAiTmYR4A9348.py"
);
confSet
(
node8Conf
);
node8
.
setConfig
(
node8Conf
);
node8
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node6"
));
System
.
out
.
println
(
"-----------------node8 节点构建成功,开始构建node9节点,依赖 node7 node8 节点--------------------"
);
PluginNode
node9
=
new
PluginNode
();
PluginNodeConfig
node9Conf
=
new
PluginNodeConfig
();
nodeSet
(
node9
,
true
);
node9
.
setName
(
"node9"
);
node9
.
setDesc
(
"node9节点"
);
node9
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa36AeryrAAAAiVRZhPM3624.py"
);
confSet
(
node9Conf
);
node9
.
setConfig
(
node9Conf
);
node9
.
setDependNodeNameList
(
Arrays
.
asList
(
"node7"
,
"node8"
));
System
.
out
.
println
(
"-----------------node9 节点构建成功,开始构建 end 节点,依赖 node9 节点--------------------"
);
PluginNode
endNode
=
new
PluginNode
();
PluginNodeConfig
endNodeConf
=
new
PluginNodeConfig
();
nodeSet
(
endNode
,
false
);
endNode
.
setName
(
"end"
);
endNode
.
setDesc
(
"end节点"
);
endNode
.
setScriptUrls
(
"ddmp/M00/00/01/CgB4Al5wa6WADyyrAAAANb42CIc7700.py"
);
confSet
(
endNodeConf
);
endNode
.
setConfig
(
endNodeConf
);
endNode
.
setDependNodeNameList
(
Collections
.
singletonList
(
"node9"
));
pluginFlow
.
setNodeList
(
Arrays
.
asList
(
startNode
,
node4
,
node5
,
node6
,
node7
,
node8
,
node9
,
endNode
));
//构建工作流信息
PluginFlowConfig
build
=
PluginFlowConfig
.
builder
().
alarmEmail
(
"huangfukexing@byitgroup.com"
)
.
alarmlAction
(
"1"
)
.
execType
(
"1"
)
.
flowCron
(
"0 0/2 * * * ? *"
)
.
flowTimeout
(
TimeUnit
.
MINUTES
.
toMillis
(
30
))
.
priority
(
"1"
)
.
repeatCount
(
1
)
.
scheduleFollow
(
"1"
)
.
build
();
pluginFlow
.
setName
(
"node4"
);
pluginFlow
.
setDesc
(
"虚拟节点"
);
pluginFlow
.
setConfig
(
build
);
pluginFlow
.
setPrincipal
(
"皇甫科星"
);
pluginFlow
.
setRePublish
(
false
);
pluginFlow
.
setAuthor
(
"huangfukexing"
);
pluginFlow
.
setType
(
"flow"
);
return
pluginFlow
;
}
/**
* 构建节点
* @param pluginNode
*/
private
static
void
nodeSet
(
PluginNode
pluginNode
,
boolean
flag
){
pluginNode
.
setAuthor
(
"皇甫"
);
pluginNode
.
setJobType
(
"SCRIPT"
);
pluginNode
.
setType
(
"node"
);
pluginNode
.
setRunCommand
(
"python ${biz_file}"
);
if
(
flag
){
ScriptParamAndPlaceholderDto
scriptParamAndPlaceholderDto
=
new
ScriptParamAndPlaceholderDto
();
Map
<
String
,
String
>
map
=
new
HashMap
<>();
map
.
put
(
"name"
,
"皇甫科星"
);
scriptParamAndPlaceholderDto
.
setPlaceholder
(
map
);
pluginNode
.
setScriptParam
(
scriptParamAndPlaceholderDto
);
}
}
/**
* 构建配置
* @param conf
*/
private
static
void
confSet
(
PluginNodeConfig
conf
){
conf
.
setFailedRetryCount
(
2
);
conf
.
setFailedRetryInterval
(
TimeUnit
.
MINUTES
.
toSeconds
(
2
));
conf
.
setNodeTimeout
(-
1L
);
conf
.
setPriority
(
"1"
);
}
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment