Skip to content

Commit ca56feb

Browse files
committed
remove powerjob relevant resource
1 parent e77fce4 commit ca56feb

7 files changed

Lines changed: 71 additions & 135 deletions

File tree

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

Lines changed: 33 additions & 86 deletions
Original file line numberDiff line numberDiff line change
@@ -28,18 +28,14 @@
2828
import com.qlangtech.tis.datax.IDataxReader;
2929
import com.qlangtech.tis.datax.IDataxWriter;
3030
import com.qlangtech.tis.datax.LifeCycleHook;
31-
import com.qlangtech.tis.datax.RpcUtils;
32-
import com.qlangtech.tis.datax.powerjob.CfgsSnapshotConsumer;
3331
import com.qlangtech.tis.datax.powerjob.ExecPhase;
3432
import com.qlangtech.tis.datax.powerjob.SplitTabSync;
3533
import com.qlangtech.tis.exec.AbstractExecContext;
3634
import com.qlangtech.tis.exec.ExecChainContextUtils;
3735
import com.qlangtech.tis.fullbuild.indexbuild.IRemoteTaskTrigger;
3836
import com.qlangtech.tis.job.common.JobParams;
3937
import com.qlangtech.tis.plugin.PluginAndCfgSnapshotLocalCache;
40-
import com.qlangtech.tis.plugin.ds.DefaultTab;
4138
import com.qlangtech.tis.plugin.ds.ISelectedTab;
42-
import com.qlangtech.tis.powerjob.SelectedTabTriggersConfig;
4339
import com.qlangtech.tis.rpc.grpc.log.ILoggerAppenderClient.LogLevel;
4440
import com.qlangtech.tis.sql.parser.tuple.creator.EntityName;
4541
import com.qlangtech.tis.web.start.TisAppLaunch;
@@ -49,7 +45,6 @@
4945
import org.apache.commons.lang.StringUtils;
5046
import org.apache.commons.lang.exception.ExceptionUtils;
5147
import org.apache.commons.lang3.tuple.Pair;
52-
import org.apache.commons.lang3.tuple.Triple;
5348
import org.slf4j.Logger;
5449
import org.slf4j.LoggerFactory;
5550

@@ -80,45 +75,39 @@ public static Pair<Boolean, JSONObject> getInstanceParams(ITaskExecutorContext c
8075

8176
public void processPostTask(ITaskExecutorContext context) {
8277

83-
RpcServiceReference svc = getRpcServiceReference();
84-
// StatusRpcClientFactory.AssembleSvcCompsite svc = statusRpc.get();
85-
Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig> pair = createExecContext(context, ExecPhase.Reduce);
86-
87-
AbstractExecContext execContext = Objects.requireNonNull(pair.getLeft(), "execContext can not be null");
88-
SelectedTabTriggersConfig triggerCfg = pair.getRight();
89-
90-
ISelectedTab tab = new DefaultTab(triggerCfg.getTabName());
91-
String postTrigger = null;
92-
Integer taskId = execContext.getTaskId();
93-
if (StringUtils.isNotEmpty(postTrigger = triggerCfg.getPostTrigger())) {
94-
95-
try {
96-
RpcUtils.setJoinStatus(taskId, false, false, svc, postTrigger);
97-
context.infoLog("exec postTrigger:{}", postTrigger);
98-
99-
IRemoteTaskTrigger postTask = createDataXJob(execContext, Pair.of(postTrigger,
100-
LifeCycleHook.Post), tab.getName());
101-
postTask.run();
102-
RpcUtils.setJoinStatus(taskId, true, false, svc, postTrigger);
103-
} catch (Exception e) {
104-
RpcUtils.setJoinStatus(taskId, true, true, svc, postTrigger);
105-
context.errorLog("postTrigger:" + postTrigger + " falid", e);
106-
throw new RuntimeException("postTrigger:" + postTrigger + " falid", e);
107-
}
108-
}
109-
110-
addSuccessPartition(context, execContext, tab.getName());
78+
// RpcServiceReference svc = getRpcServiceReference();
79+
// // StatusRpcClientFactory.AssembleSvcCompsite svc = statusRpc.get();
80+
// Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig> pair = createExecContext(context, ExecPhase.Reduce);
81+
//
82+
// AbstractExecContext execContext = Objects.requireNonNull(pair.getLeft(), "execContext can not be null");
83+
// SelectedTabTriggersConfig triggerCfg = pair.getRight();
84+
//
85+
// ISelectedTab tab = new DefaultTab(triggerCfg.getTabName());
86+
// String postTrigger = null;
87+
// Integer taskId = execContext.getTaskId();
88+
// if (StringUtils.isNotEmpty(postTrigger = triggerCfg.getPostTrigger())) {
89+
//
90+
// try {
91+
// RpcUtils.setJoinStatus(taskId, false, false, svc, postTrigger);
92+
// context.infoLog("exec postTrigger:{}", postTrigger);
93+
//
94+
// IRemoteTaskTrigger postTask = createDataXJob(execContext, Pair.of(postTrigger,
95+
// LifeCycleHook.Post), tab.getName());
96+
// postTask.run();
97+
// RpcUtils.setJoinStatus(taskId, true, false, svc, postTrigger);
98+
// } catch (Exception e) {
99+
// RpcUtils.setJoinStatus(taskId, true, true, svc, postTrigger);
100+
// context.errorLog("postTrigger:" + postTrigger + " falid", e);
101+
// throw new RuntimeException("postTrigger:" + postTrigger + " falid", e);
102+
// }
103+
// }
104+
//
105+
// addSuccessPartition(context, execContext, tab.getName());
111106
}
112107

113108

114-
/**
115-
* initialize 节点之后执行的任务节点
116-
*
117-
* @param context
118-
* @return //@throws com.qlangtech.tis.datax.powerjob.InstanceParamsException
119-
*/
120-
public static Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig>
121-
createExecContext(ITaskExecutorContext context, ExecPhase execPhase) {
109+
// public static Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig>
110+
// createExecContext(ITaskExecutorContext context, ExecPhase execPhase) {
122111
// JSONObject instanceParams = null;
123112
// Pair<Boolean, JSONObject> instanceParamsGetter = getInstanceParams(context);
124113
// instanceParams = instanceParamsGetter.getRight();
@@ -136,52 +125,10 @@ public void processPostTask(ITaskExecutorContext context) {
136125
// + triggerCfg.getSplitTabsCfg().stream().map(CuratorDataXTaskMessage::getJobName).collect(Collectors.joining(",")));
137126
//
138127
// return pair;
139-
throw new UnsupportedOperationException();
140-
}
128+
// throw new UnsupportedOperationException();
129+
// }
130+
141131

142-
private static Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig> //
143-
createExecContext(ITaskExecutorContext context, Integer taskId, JSONObject instanceParams) {
144-
throw new UnsupportedOperationException();
145-
// if (taskId == null) {
146-
// throw new IllegalArgumentException("param taskId can not be null");
147-
// }
148-
//
149-
// SelectedTabTriggersConfig triggerCfg = getTriggerCfg(context);
150-
//
151-
// final CfgsSnapshotConsumer snapshotConsumer = new CfgsSnapshotConsumer();
152-
// AbstractExecContext execContext = IExecChainContext.deserializeInstanceParams(triggerCfg, instanceParams, (ctx) -> {
153-
// ctx.setLatestPhaseStatusCollection(cacheSnaphsot.getPreviousStatus(ctx.getTaskId(), () -> {
154-
// Integer prevTaskId = instanceParams.getInteger(JobParams.KEY_PREVIOUS_TASK_ID);
155-
// if (prevTaskId == null) {
156-
// return null;
157-
// }
158-
// return getRpcServiceReference().loadPhaseStatusFromLatest(prevTaskId);
159-
// }));
160-
//
161-
// }, snapshotConsumer);
162-
//
163-
// execContext.setSpecifiedLocalLoggerPath(context.getSpecifiedLocalLoggerPath());
164-
// execContext.setDisableGrpcRemoteServerConnect(context.isDisableGrpcRemoteServerConnect());
165-
// // execContext.setAttribute(JobCommon.KEY_TASK_ID, Objects.requireNonNull(taskId, "taskId can not be null"));
166-
//
167-
//
168-
// /**
169-
// * 同步必要的配置及tpi资源到本地
170-
// */
171-
// snapshotConsumer.synchronizTpisAndConfs(execContext, cacheSnaphsot);
172-
//
173-
// Long triggerTimestamp = execContext.getPartitionTimestampWithMillis();// instanceParams.getLong(DataxUtils
174-
// System.setProperty(DataxUtils.EXEC_TIMESTAMP, String.valueOf(triggerTimestamp));
175-
// for (CuratorDataXTaskMessage tskMsg : triggerCfg.getSplitTabsCfg()) {
176-
// tskMsg.setExecTimeStamp(triggerTimestamp);
177-
// }
178-
// triggerCfg.getSplitTabsCfg().forEach((tskMsg) -> {
179-
// tskMsg.setJobId(taskId);
180-
// });
181-
//
182-
//
183-
// return Triple.of(execContext, snapshotConsumer, triggerCfg);
184-
}
185132

186133
public static DataxPrePostConsumer createPrePostConsumer() {
187134
DataXJobRunEnvironmentParamsSetter.ExtraJavaSystemPramsSuppiler systemPramsSuppiler = createSysPramsSuppiler();

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020

2121
import com.alibaba.fastjson.JSONObject;
2222
import com.qlangtech.tis.cloud.ITISCoordinator;
23+
import com.qlangtech.tis.datax.DataXName;
2324
import com.qlangtech.tis.datax.IDataxProcessor;
2425
import com.qlangtech.tis.datax.IDataxWriter;
2526
import com.qlangtech.tis.datax.RpcUtils;
@@ -32,7 +33,6 @@
3233
import com.qlangtech.tis.exec.ExecChainContextUtils;
3334
import com.qlangtech.tis.fullbuild.indexbuild.IPartionableWarehouse;
3435
import com.qlangtech.tis.job.common.JobParams;
35-
import com.qlangtech.tis.powerjob.TriggersConfig;
3636
import com.qlangtech.tis.rpc.grpc.log.ILoggerAppenderClient.LogLevel;
3737
import com.qlangtech.tis.sql.parser.ISqlTask;
3838
import com.qlangtech.tis.sql.parser.SqlTaskNodeMeta;
@@ -112,7 +112,7 @@ protected void process(ITaskExecutorContext context) throws Exception {
112112
private AbstractExecContext createDftExecContent(ITaskExecutorContext context) {
113113
JSONObject instanceParams = (context.getInstanceParams());
114114
final CfgsSnapshotConsumer snapshotConsumer = new CfgsSnapshotConsumer();
115-
TriggersConfig triggerCfg = new TriggersConfig(instanceParams.getString(JobParams.KEY_COLLECTION), StoreResourceType.DataFlow);
115+
DataXName triggerCfg = new DataXName(instanceParams.getString(JobParams.KEY_COLLECTION), StoreResourceType.DataFlow);
116116
// FIXME shall initialize execContext
117117
AbstractExecContext execContext = null;
118118
// IExecChainContext.deserializeInstanceParams(triggerCfg, instanceParams, (ctx) -> {

tis-datax/executor/tis-datax-executor/src/main/java/com/qlangtech/tis/datax/join/DataXJoinProcessExecutor.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,6 @@
66
import com.qlangtech.tis.datax.IDataxProcessor;
77
import com.qlangtech.tis.datax.IDataxWriter;
88
import com.qlangtech.tis.datax.RpcUtils;
9-
import com.qlangtech.tis.datax.StoreResourceType;
109
import com.qlangtech.tis.exec.AbstractExecContext;
1110
import com.qlangtech.tis.exec.ExecutePhaseRange;
1211
import com.qlangtech.tis.exec.ExecuteResult;
@@ -20,7 +19,6 @@
2019
import com.qlangtech.tis.job.common.JobParams;
2120
import com.qlangtech.tis.offline.DataxUtils;
2221
import com.qlangtech.tis.plugin.ds.IDataSourceFactoryGetter;
23-
import com.qlangtech.tis.powerjob.TriggersConfig;
2422
import com.qlangtech.tis.sql.parser.SqlTaskNodeMeta;
2523
import com.qlangtech.tis.sql.parser.TopologyDir;
2624
import com.qlangtech.tis.sql.parser.er.IPrimaryTabFinder;
@@ -194,7 +192,7 @@ public static JSONObject deserializeInstanceParams(CommandLine line) {
194192
private static AbstractExecContext createDftExecContent(CommandLine line) {
195193
JSONObject instanceParams = deserializeInstanceParams(line);
196194

197-
TriggersConfig triggersConfig = new TriggersConfig(instanceParams.getString(JobParams.KEY_COLLECTION), StoreResourceType.DataFlow);
195+
// TriggersConfig triggersConfig = new TriggersConfig(instanceParams.getString(JobParams.KEY_COLLECTION), StoreResourceType.DataFlow);
198196
// FIXME shall initialize execContext instance
199197
AbstractExecContext execContext = null; //IExecChainContext.deserializeInstanceParams(triggersConfig, instanceParams);
200198
// execContext.setResType(StoreResourceType.DataFlow);

tis-datax/executor/tis-datax-executor/src/test/java/com/qlangtech/tis/datax/executor/BasicTISTableDumpProcessorTest.java

Lines changed: 4 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,6 @@
3434
import com.qlangtech.tis.datax.impl.DataxProcessor;
3535
import com.qlangtech.tis.datax.impl.DataxWriter.BaseDataxWriterDescriptor;
3636
import com.qlangtech.tis.datax.impl.TransformerInfo;
37-
import com.qlangtech.tis.datax.powerjob.CfgsSnapshotConsumer;
38-
import com.qlangtech.tis.datax.powerjob.ExecPhase;
39-
import com.qlangtech.tis.exec.AbstractExecContext;
4037
import com.qlangtech.tis.exec.IExecChainContext;
4138
import com.qlangtech.tis.job.common.JobParams;
4239
import com.qlangtech.tis.manage.biz.dal.pojo.Application;
@@ -48,13 +45,10 @@
4845
import com.qlangtech.tis.plugin.trigger.JobTrigger;
4946
import com.qlangtech.tis.powerjob.SelectedTabTriggers;
5047
import com.qlangtech.tis.powerjob.SelectedTabTriggers.PowerJobRemoteTaskTrigger;
51-
import com.qlangtech.tis.powerjob.SelectedTabTriggersConfig;
5248
import com.qlangtech.tis.test.TISEasyMock;
5349
import com.qlangtech.tis.util.IPluginContext;
5450
import org.apache.commons.lang3.tuple.Pair;
55-
import org.apache.commons.lang3.tuple.Triple;
5651
import org.easymock.EasyMock;
57-
import org.junit.Assert;
5852
import org.junit.Before;
5953
import org.junit.Test;
6054

@@ -121,11 +115,11 @@ public void createExecContext() {
121115
instanceParams.put(JobParams.KEY_TASK_ID, taskId);
122116

123117
replay();
124-
Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig> result
125-
= BasicTISTableDumpProcessor.createExecContext(context, ExecPhase.Reduce);
118+
// Triple<AbstractExecContext, CfgsSnapshotConsumer, SelectedTabTriggersConfig> result
119+
// = BasicTISTableDumpProcessor.createExecContext(context, ExecPhase.Reduce);
126120

127-
AbstractExecContext execContext = result.getLeft();
128-
Assert.assertEquals(taskId, (Integer) execContext.getTaskId());
121+
// AbstractExecContext execContext = result.getLeft();
122+
// Assert.assertEquals(taskId, (Integer) execContext.getTaskId());
129123

130124
verifyAll();
131125
}

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

Lines changed: 22 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,9 @@
1818

1919
package com.qlangtech.tis.plugin.datax;
2020

21-
import com.alibaba.fastjson.JSONObject;
2221
import com.qlangtech.tis.assemble.ExecResult;
2322
import com.qlangtech.tis.dao.ICommonDAOContext;
2423
import com.qlangtech.tis.datax.IDataxProcessor;
25-
import com.qlangtech.tis.datax.job.ITISPowerJob;
26-
import com.qlangtech.tis.trigger.util.JsonUtil;
2724
import com.qlangtech.tis.workflow.pojo.WorkFlowBuildHistory;
2825
import com.qlangtech.tis.workflow.pojo.WorkFlowBuildHistoryCriteria;
2926
import org.slf4j.Logger;
@@ -58,29 +55,30 @@ public Integer getTisTaskId() {
5855

5956
public Long getSPIWorkflowInstanceId() {
6057

61-
if (this.spiWorkflowInstanceId == null) {
62-
WorkFlowBuildHistory wfBuildHistory
63-
= daoContext.getTaskBuildHistoryDAO().selectByPrimaryKey(tisTaskId);
64-
this.spiWorkflowInstanceId = ITISPowerJob.getPowerJobWorkflowInstanceId(wfBuildHistory, true);
65-
}
66-
return this.spiWorkflowInstanceId;
58+
throw new UnsupportedOperationException();
59+
// if (this.spiWorkflowInstanceId == null) {
60+
// WorkFlowBuildHistory wfBuildHistory
61+
// = daoContext.getTaskBuildHistoryDAO().selectByPrimaryKey(tisTaskId);
62+
// this.spiWorkflowInstanceId = ITISPowerJob.getPowerJobWorkflowInstanceId(wfBuildHistory, true);
63+
// }
64+
// return this.spiWorkflowInstanceId;
6765
}
6866

69-
public void setSPIWorkflowInstanceId(Long workflowInstanceId) {
70-
// logger.info("create workflow instanceId:{}", workflowInstanceId);
71-
// 需要将task执行历史记录更新,将instanceId 绑定到历史记录上去,以便后续最终
72-
WorkFlowBuildHistory record = new WorkFlowBuildHistory();
73-
JSONObject wfHistory = new JSONObject();
74-
wfHistory.put(ITISPowerJob.KEY_POWERJOB_WORKFLOW_INSTANCE_ID, workflowInstanceId);
75-
record.setAsynSubTaskStatus(JsonUtil.toString(wfHistory));
76-
WorkFlowBuildHistoryCriteria taskHistoryCriteria = new WorkFlowBuildHistoryCriteria();
77-
taskHistoryCriteria.createCriteria().andIdEqualTo(tisTaskId);
78-
if (daoContext.getTaskBuildHistoryDAO().updateByExampleSelective(record, taskHistoryCriteria) < 1) {
79-
throw new IllegalStateException("update taskBuildHistory faild,taskId:" + tisTaskId
80-
+ ",powerJob workflowInstanceId:" + workflowInstanceId);
81-
}
82-
this.spiWorkflowInstanceId = workflowInstanceId;
83-
}
67+
// public void setSPIWorkflowInstanceId(Long workflowInstanceId) {
68+
// // logger.info("create workflow instanceId:{}", workflowInstanceId);
69+
// // 需要将task执行历史记录更新,将instanceId 绑定到历史记录上去,以便后续最终
70+
// WorkFlowBuildHistory record = new WorkFlowBuildHistory();
71+
// JSONObject wfHistory = new JSONObject();
72+
// wfHistory.put(ITISPowerJob.KEY_POWERJOB_WORKFLOW_INSTANCE_ID, workflowInstanceId);
73+
// record.setAsynSubTaskStatus(JsonUtil.toString(wfHistory));
74+
// WorkFlowBuildHistoryCriteria taskHistoryCriteria = new WorkFlowBuildHistoryCriteria();
75+
// taskHistoryCriteria.createCriteria().andIdEqualTo(tisTaskId);
76+
// if (daoContext.getTaskBuildHistoryDAO().updateByExampleSelective(record, taskHistoryCriteria) < 1) {
77+
// throw new IllegalStateException("update taskBuildHistory faild,taskId:" + tisTaskId
78+
// + ",powerJob workflowInstanceId:" + workflowInstanceId);
79+
// }
80+
// this.spiWorkflowInstanceId = workflowInstanceId;
81+
// }
8482

8583
/**
8684
*

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@
2020
import com.qlangtech.tis.exec.AbstractExecContext;
2121
import com.qlangtech.tis.fullbuild.indexbuild.IRemoteTaskTrigger;
2222
import com.qlangtech.tis.powerjob.SelectedTabTriggers;
23-
import com.qlangtech.tis.powerjob.TriggersConfig;
23+
//import com.qlangtech.tis.powerjob.TriggersConfig;
2424
import com.qlangtech.tis.powerjob.model.InstanceStatus;
2525
import com.qlangtech.tis.powerjob.model.PEWorkflowDAG;
2626
import com.tis.hadoop.rpc.RpcServiceReference;
@@ -250,7 +250,7 @@ private void executeTaskInternal(TaskExecutionMessage msg) throws Exception {
250250
logger.info("Executing DataX job: nodeName={}, workflowContext={}",
251251
node.getNodeName(), msg.getWorkflowContext());
252252
AbstractExecContext execContext = TaskExecutionMessage.deserializeInstanceParams(
253-
TriggersConfig.create(msg.getDataXName()) //
253+
(msg.getDataXName()) //
254254
, msg //
255255
, (ctx) -> {
256256
} //

0 commit comments

Comments
 (0)