Skip to content

Commit d36e619

Browse files
committed
add akka dag support ,issue:datavane/tis#486
1 parent a146f8e commit d36e619

31 files changed

Lines changed: 770 additions & 250 deletions

File tree

pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@
6666
<module>tis-powerjob-common-plugin</module>
6767

6868
<module>tis-transformer</module>
69-
<module>tis-datax/tis-datax-dolphinscheduler-plugin</module>
69+
<!-- <module>tis-datax/tis-datax-dolphinscheduler-plugin</module>-->
7070
<module>tis-datax/tis-hive-shim-common</module>
7171

7272

tis-datax/.claude/settings.local.json

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,10 @@
44
"Bash(mvn test -pl tis-datax-local-akka-executor -Dtest=com.qlangtech.tis.dag.actor.TestWorkflowInstanceActor)",
55
"Bash(find /Users/mozhenghua/j2ee_solution/project -type f -name \"*.java\" -exec grep -l \"class WorkflowDAGFileManager\" {} ;)",
66
"Bash(find /Users/mozhenghua/j2ee_solution/project -type f -name \"*.java\" -exec grep -l \"enum InstanceStatus\" {} ;)",
7-
"Bash(find /Users/mozhenghua/j2ee_solution/project -type f -name \"*.java\" -exec grep -l \"enum WorkflowNodeType\" {} ;)"
7+
"Bash(find /Users/mozhenghua/j2ee_solution/project -type f -name \"*.java\" -exec grep -l \"enum WorkflowNodeType\" {} ;)",
8+
"Bash(mvn compile -Dmaven.test.skip=true)",
9+
"Bash(mvn install -Dmaven.test.skip=true)",
10+
"Bash(mvn compile -pl tis-datax-local-akka-executor -am -Dmaven.test.skip=true -q)"
811
]
912
}
1013
}

tis-datax/executor/tis-datax-executor/src/main/java/com/qlangtech/tis/datax/DataXJobSingleProcessorExecutor.java

Lines changed: 13 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818

1919
package com.qlangtech.tis.datax;
2020

21+
import com.qlangtech.tis.job.common.JobCommon;
2122
import com.qlangtech.tis.manage.common.Config;
2223
import com.qlangtech.tis.manage.common.TisUTF8;
2324
import com.qlangtech.tis.offline.DataxUtils;
@@ -44,24 +45,22 @@
4445
**/
4546
public abstract class DataXJobSingleProcessorExecutor<T extends IDataXTaskRelevant> {
4647
private static final Logger logger = LoggerFactory.getLogger(DataXJobSingleProcessorExecutor.class);
47-
48+
public static final int DEFAULT_TASK_ID = 999;
4849
// 记录当前正在执行的任务<taskid,ExecuteWatchdog>
4950
public final ConcurrentHashMap<Integer, ExecuteWatchdog> runningTask = new ConcurrentHashMap<>();
5051

5152
public void consumeMessage(T msg) throws Exception {
52-
//MDC.put();
53-
throw new UnsupportedOperationException();
54-
// Integer taskId = msg.getTaskId();
55-
// String jobName = msg.getJobName();
56-
// String dataxName = msg.getDataXName();
57-
// // StoreResourceType resType = Objects.requireNonNull(msg.getResType(), "resType can not be null");
58-
// // MDC.put(JobCommon.KEY_TASK_ID, String.valueOf(jobId));
59-
// // MDC.put(JobCommon.KEY_COLLECTION, dataxName);
60-
// JobCommon.setMDC(taskId, dataxName);
61-
//
62-
//
63-
// // 查看当前任务是否正在进行中,如果已经终止则要退出
64-
// execSystemTask(msg, taskId, jobName, dataxName);
53+
// Integer taskId = PreviewLaunchParam;// msg.getTaskId();
54+
String jobName = msg.getJobName();
55+
DataXName dataxName = msg.getDataXName();
56+
// StoreResourceType resType = Objects.requireNonNull(msg.getResType(), "resType can not be null");
57+
// MDC.put(JobCommon.KEY_TASK_ID, String.valueOf(jobId));
58+
// MDC.put(JobCommon.KEY_COLLECTION, dataxName);
59+
JobCommon.setMDC(DEFAULT_TASK_ID, dataxName.getPipelineName());
60+
61+
62+
// 查看当前任务是否正在进行中,如果已经终止则要退出
63+
execSystemTask(msg, DEFAULT_TASK_ID, jobName, dataxName.getPipelineName());
6564
}
6665

6766
protected void execSystemTask(T msg, Integer jobId, String jobName, String dataxName) throws IOException,

tis-datax/executor/tis-datax-executor/src/main/java/com/qlangtech/tis/datax/powerjob/SplitTabSync.java

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -32,16 +32,15 @@ public SplitTabSync(DataXJobInfo tskMsg, Integer allRows) {
3232
// this.taskConfig = Objects.requireNonNull(taskConfig);
3333
}
3434

35-
public void execSync(final AbstractExecContext execChainContext, RpcServiceReference statusRpc) {
35+
public IRemoteTaskTrigger createTrigger(final AbstractExecContext execChainContext, RpcServiceReference statusRpc) {
3636
if (statusRpc == null) {
3737
throw new IllegalArgumentException("statusRpc can not be null");
3838
}
3939
DataXJobSubmit dataXJobSubmit = getDataXJobSubmit(execChainContext);
4040
if (dataXJobSubmit instanceof DataXJobRunEnvironmentParamsSetter) {
4141
DataXJobRunEnvironmentParamsSetter runEnvironmentParamsSetter =
4242
(DataXJobRunEnvironmentParamsSetter) dataXJobSubmit;
43-
DataxPrePostConsumer prePostConsumer = BasicTISTableDumpProcessor.createPrePostConsumer();// new
44-
// DataxPrePostConsumer();
43+
DataxPrePostConsumer prePostConsumer = BasicTISTableDumpProcessor.createPrePostConsumer();
4544
runEnvironmentParamsSetter.setClasspath(prePostConsumer.getClasspath());
4645
runEnvironmentParamsSetter.setWorkingDirectory(prePostConsumer.getWorkingDirectory());
4746
runEnvironmentParamsSetter.setExtraJavaSystemPramsSuppiler(prePostConsumer.getExtraJavaSystemPramsSuppiler());
@@ -50,10 +49,11 @@ public void execSync(final AbstractExecContext execChainContext, RpcServiceRefer
5049
DataXJobSubmit.IDataXJobContext dataXJobContext = DataXJobSubmit.IDataXJobContext.create(execChainContext);
5150
IDataxProcessor processor = execChainContext.getProcessor();
5251
CuratorDataXTaskMessage taskCfg = dataXJobSubmit.getDataXJobDTO(dataXJobContext, tskMsg, processor, this.allRows);
53-
IRemoteTaskTrigger tskTrigger = dataXJobSubmit.createDataXJob(dataXJobContext, statusRpc,
54-
tskMsg, processor, taskCfg);
52+
return dataXJobSubmit.createDataXJob(dataXJobContext, statusRpc, tskMsg, processor, taskCfg);
53+
}
5554

56-
tskTrigger.run();
55+
public void execSync(final AbstractExecContext execChainContext, RpcServiceReference statusRpc) {
56+
createTrigger(execChainContext, statusRpc).run();
5757
}
5858

5959
private static DataXJobSubmit getDataXJobSubmit(AbstractExecContext execChainContext) {

tis-datax/pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,7 @@
4141

4242
</properties>
4343
<modules>
44-
<module>tis-datax-local-powerjob-executor</module>
44+
<!-- <module>tis-datax-local-powerjob-executor</module>-->
4545
<module>tis-aliyun-jindo-sdk-extends</module>
4646
<module>tis-datax-common-plugin</module>
4747
<module>tis-datax-common-rdbms-plugin</module>
@@ -90,8 +90,8 @@
9090

9191

9292
<!-- <module>tis-datax-rabbitmq-plugin</module>-->
93-
<module>executor/powerjob-worker-samples</module>
94-
<module>executor/dolphinscheduler-task-tis-datasync</module>
93+
<!-- <module>executor/powerjob-worker-samples</module>-->
94+
<!-- <module>executor/dolphinscheduler-task-tis-datasync</module>-->
9595
<module>executor/tis-datax-executor</module>
9696
<module>tis-datax-local-executor-utils</module>
9797
<module>tis-datax-kingbase-plugin</module>

tis-datax/tis-datax-common-plugin/src/main/java/com/qlangtech/tis/plugin/datax/BasicDistributedSPIDataXJobSubmit.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -102,13 +102,13 @@ public TriggerBuildResult triggerWorkflowJob(
102102
* @return
103103
*/
104104
@Override
105-
public TriggerBuildResult triggerJob(IExecChainContext execChainContext, DataXName appName) {
105+
public final TriggerBuildResult triggerJob(IExecChainContext execChainContext, DataXName appName) {
106106
if ((appName) == null) {
107107
throw new IllegalArgumentException("appName " + appName + " can not be empty");
108108
}
109109

110110
BasicWorkflowPayload<WF_INSTANCE> appPayload = createApplicationPayload(execChainContext, appName);
111-
return appPayload.triggerWorkflow(execChainContext,getStatusRpc());
111+
return appPayload.triggerWorkflow(execChainContext, getStatusRpc());
112112
}
113113

114114
protected abstract BasicWorkflowPayload<WF_INSTANCE> createWorkflowPayload(

tis-datax/tis-datax-dolphinscheduler-plugin/src/main/java/com/qlangtech/tis/plugin/datax/doplinscheduler/DSWorkflowPayload.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -337,7 +337,8 @@ private void innerSaveJob(
337337
//======================================
338338
tabTriggers = entry.getValue();
339339
selectedTab = entry.getKey();
340-
JSONObject mrParams = new JSONObject(tabTriggers.createMRParams());
340+
// tabTriggers.createMRParams()
341+
JSONObject mrParams = new JSONObject();
341342
this.setDisableGrpcRemoteServerConnect(mrParams);
342343
if (tabTriggers.getPostTrigger() != null) {
343344
containPostTrigger = true;

tis-datax/tis-datax-dolphinscheduler-plugin/src/main/java/com/qlangtech/tis/plugin/datax/doplinscheduler/DolphinschedulerDistributedSPIDataXJobSubmit.java

Lines changed: 8 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,6 @@
2121
import com.alibaba.citrus.turbine.Context;
2222
import com.qlangtech.tis.annotation.Public;
2323
import com.qlangtech.tis.build.task.IBuildHistory;
24-
import com.qlangtech.tis.coredefine.module.action.TriggerBuildResult;
2524
import com.qlangtech.tis.dao.ICommonDAOContext;
2625
import com.qlangtech.tis.datax.DataXName;
2726
import com.qlangtech.tis.datax.DefaultDataXProcessorManipulate;
@@ -64,14 +63,14 @@ public boolean cancelTask(IControlMsgHandler module, Context context, IBuildHist
6463
throw TisException.create("Dolphinscheduler DAG job cancel is not supported");
6564
}
6665

67-
@Override
68-
public TriggerBuildResult triggerJob(IExecChainContext execChainContext, DataXName appName //Optional<Long> workflowInstanceIdOpt,
69-
// , Optional<PhaseStatusCollection> latestSuccessWorkflowHistory
70-
) {
71-
return super.triggerJob(execChainContext, appName //workflowInstanceIdOpt,
72-
// , latestSuccessWorkflowHistory
73-
);
74-
}
66+
// @Override
67+
// public TriggerBuildResult triggerJob(IExecChainContext execChainContext, DataXName appName //Optional<Long> workflowInstanceIdOpt,
68+
// // , Optional<PhaseStatusCollection> latestSuccessWorkflowHistory
69+
// ) {
70+
// return super.triggerJob(execChainContext, appName //workflowInstanceIdOpt,
71+
// // , latestSuccessWorkflowHistory
72+
// );
73+
// }
7574

7675
@Override
7776
protected BasicWorkflowPayload createWorkflowPayload(

tis-datax/tis-datax-dolphinscheduler-plugin/src/main/java/com/qlangtech/tis/plugin/datax/doplinscheduler/history/DSWorkFlowBuildHistoryPayload.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@
2424
import com.qlangtech.tis.datax.DefaultDataXProcessorManipulate;
2525
import com.qlangtech.tis.datax.IDataxProcessor;
2626
import com.qlangtech.tis.extension.IDescribableManipulate;
27-
import com.qlangtech.tis.manage.IAppSource;
2827
import com.qlangtech.tis.plugin.datax.WorkFlowBuildHistoryPayload;
2928
import com.qlangtech.tis.plugin.datax.doplinscheduler.export.DolphinSchedulerURLBuilder.DolphinSchedulerResponse;
3029
import com.qlangtech.tis.plugin.datax.doplinscheduler.export.ExportTISPipelineToDolphinscheduler;
@@ -44,7 +43,7 @@ public DSWorkFlowBuildHistoryPayload(IDataxProcessor dataxProcessor, Integer tis
4443
super(dataxProcessor, tisTaskId, daoContext);
4544
//Pair<List<ExportTISPipelineToDolphinscheduler>, IPluginStore<DefaultDataXProcessorManipulate>>
4645
DefaultDataXProcessorManipulate.ProcessorManipulateManager<ExportTISPipelineToDolphinscheduler> pluginStorePair
47-
= DefaultDataXProcessorManipulate.loadPlugins(null, ExportTISPipelineToDolphinscheduler.class, ((IAppSource) this.dataxProcessor).getDataXName()
46+
= DefaultDataXProcessorManipulate.loadPlugins(null, ExportTISPipelineToDolphinscheduler.class, ( this.dataxProcessor).getDataXName()
4847
, new IDescribableManipulate.IManipulateStorable() {
4948
@Override
5049
public boolean isManipulateStorable() {

tis-datax/tis-datax-local-akka-executor/src/main/java/com/qlangtech/tis/dag/BatchJobCrontab.java

Lines changed: 27 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,15 @@
22

33
import com.alibaba.citrus.turbine.Context;
44
import com.alibaba.fastjson.JSONObject;
5+
import com.google.common.collect.Lists;
6+
import com.qlangtech.tis.assemble.TriggerType;
57
import com.qlangtech.tis.datax.DefaultDataXProcessorManipulate;
68
import com.qlangtech.tis.datax.IManipulateStatus;
79
import com.qlangtech.tis.datax.TimeFormat;
810
import com.qlangtech.tis.extension.Descriptor;
911
import com.qlangtech.tis.extension.TISExtension;
12+
import com.qlangtech.tis.manage.common.Config;
13+
import com.qlangtech.tis.manage.common.HttpUtils;
1014
import com.qlangtech.tis.plugin.IdentityDesc;
1115
import com.qlangtech.tis.plugin.annotation.FormField;
1216
import com.qlangtech.tis.plugin.annotation.FormFieldType;
@@ -15,9 +19,11 @@
1519
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
1620
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
1721
import com.qlangtech.tis.util.IPluginContext;
22+
import org.apache.commons.lang3.StringUtils;
1823

1924
import java.util.Collections;
2025
import java.util.Date;
26+
import java.util.List;
2127
import java.util.Optional;
2228

2329
/**
@@ -44,9 +50,29 @@ public JSONObject describePlugin() {
4450
return Descriptor.getManipulateMeta(false, this);
4551
}
4652

53+
/**
54+
* //@see DAGWorkflowServlet
55+
*
56+
* @param pluginContext
57+
* @param context
58+
* @param itemsProcessor
59+
*/
4760
@Override
4861
protected void afterManipuldateProcess(IPluginContext pluginContext, Optional<Context> context, ManipulateItemsProcessor itemsProcessor) {
4962

63+
64+
final String getAssembleHttpHost = Config.getAssembleHttpHost();
65+
if (StringUtils.isEmpty(getAssembleHttpHost)) {
66+
throw new IllegalArgumentException("param getAssembleHttpHost can not be empty");
67+
}
68+
69+
List<HttpUtils.PostParam> params = Lists.newArrayList();
70+
// params.add(new HttpUtils.PostParam(HttpUtils.KEY_METHOD, HttpUtils.KEY_METHOD_HANDLE_REGISTER_SCHEDULE));
71+
params.addAll(HttpUtils.dataXToParams(pluginContext.getCollectionName()));
72+
params.add(new HttpUtils.PostParam(TriggerType.KEY_CONTAB, this.crontab));
73+
params.add(new HttpUtils.PostParam(TriggerType.KEY_CRONTAB_TURN_ON, !itemsProcessor.isDeleteProcess() && this.turnOn));
74+
HttpUtils.postAssembleDAGServlet(HttpUtils.KEY_METHOD_HANDLE_REGISTER_SCHEDULE, params, (respStream) -> null);
75+
5076
}
5177

5278
@Override
@@ -86,7 +112,7 @@ protected boolean verify(IControlMsgHandler msgHandler, Context context, PostFor
86112
// CrontabTriggerStrategy describable = postFormVals.newInstance();
87113
Date nextFireTime = getNextFireTime(msgHandler, context, KEY_CRONTAB, postFormVals.getField(KEY_CRONTAB));
88114
if (nextFireTime != null) {
89-
msgHandler.addActionMessage(context, "最近触发时间为:" + TimeFormat.yyyyMMddHHmmss.format(nextFireTime.getTime()));
115+
msgHandler.addActionMessage(context, "最近触发时间为:" + TimeFormat.yyyyMMdd_HH_mm_ss.format(nextFireTime.getTime()));
90116
return true;
91117
}
92118
return false;

0 commit comments

Comments
 (0)