Skip to content

Commit e77fce4

Browse files
committed
add DefaultAKKADataXWorkerLauncher
1 parent 5cbf6db commit e77fce4

8 files changed

Lines changed: 76 additions & 72 deletions

File tree

tis-datax/tis-aliyun-jindo-sdk-extends/tis-aliyun-jindo-sdk-extends-impl/dependency-reduced-pom.xml

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
<parent>
44
<artifactId>tis-aliyun-jindo-sdk-extends</artifactId>
55
<groupId>com.qlangtech.tis.plugins</groupId>
6-
<version>5.0.0</version>
6+
<version>5.1.0</version>
77
</parent>
88
<modelVersion>4.0.0</modelVersion>
99
<groupId>com.qlangtech.tis.plugins</groupId>
1010
<artifactId>tis-aliyun-jindo-sdk-extends-impl</artifactId>
11-
<version>5.0.0</version>
11+
<version>5.1.0</version>
1212
<licenses>
1313
<license>
1414
<name>GNU Affero General Public License</name>
@@ -46,7 +46,7 @@
4646
<dependency>
4747
<groupId>com.alibaba.datax</groupId>
4848
<artifactId>datax-common</artifactId>
49-
<version>5.0.0</version>
49+
<version>5.1.0</version>
5050
<scope>provided</scope>
5151
</dependency>
5252
<dependency>
@@ -64,7 +64,7 @@
6464
<dependency>
6565
<groupId>com.qlangtech.tis</groupId>
6666
<artifactId>tis-plugin</artifactId>
67-
<version>5.0.0</version>
67+
<version>5.1.0</version>
6868
<scope>provided</scope>
6969
</dependency>
7070
<dependency>

tis-datax/tis-datax-local-akka-executor/src/main/java/com/qlangtech/tis/DataXWorkerLauncher.java renamed to tis-datax/tis-datax-local-akka-executor/src/main/java/com/qlangtech/tis/DefaultAKKADataXWorkerLauncher.java

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,9 @@
22

33
import com.qlangtech.tis.dag.TISActorSystem;
44
import com.qlangtech.tis.datax.DataXJobSubmitParams;
5+
import com.qlangtech.tis.extension.TISExtension;
56
import com.qlangtech.tis.plugin.akka.DAORestDelegateFacade;
7+
import com.qlangtech.tis.plugin.datax.DataXWorkerLauncher;
68
import org.slf4j.Logger;
79
import org.slf4j.LoggerFactory;
810

@@ -14,13 +16,16 @@
1416
* export AKKA_PORT=2552
1517
* export AKKA_SEED_NODES="akka://TIS-DAG-System@192.168.28.189:2551"
1618
* </pre>
19+
*
1720
* @author 百岁 (baisui@qlangtech.com)
1821
* @date 2026/3/5
1922
*/
20-
public class DataXWorkerLauncher {
21-
private static final Logger logger = LoggerFactory.getLogger(DataXWorkerLauncher.class);
23+
@TISExtension
24+
public class DefaultAKKADataXWorkerLauncher extends DataXWorkerLauncher {
25+
private static final Logger logger = LoggerFactory.getLogger(DefaultAKKADataXWorkerLauncher.class);
2226

23-
public static void main(String[] args) {
27+
@Override
28+
public void start() {
2429
DataXJobSubmitParams submitParams = DataXJobSubmitParams.getDftIfEmpty();
2530
logger.info("start to launch DataX worker,maxInstancesPerNode:{},maxTotalNrOfInstances:{}", submitParams.maxInstancesPerNode, submitParams.maxTotalNrOfInstances);
2631
DAORestDelegateFacade akkaClusterDependenceDao = DAORestDelegateFacade.createAKKAClusterDependenceDao();
@@ -31,4 +36,9 @@ public static void main(String[] args) {
3136
tisActorSystem.initialize();
3237
logger.info("launch DataX Worker successful");
3338
}
39+
40+
public static void main(String[] args) {
41+
DefaultAKKADataXWorkerLauncher launcher = new DefaultAKKADataXWorkerLauncher();
42+
launcher.start();
43+
}
3444
}

tis-incr/tis-flink-cdc-kingbase-shade-4-debezium-connector-postgresql/dependency-reduced-pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
<parent>
44
<artifactId>tis-incr</artifactId>
55
<groupId>com.qlangtech.tis.plugins</groupId>
6-
<version>5.0.0</version>
6+
<version>5.1.0</version>
77
</parent>
88
<modelVersion>4.0.0</modelVersion>
99
<groupId>com.qlangtech.tis.plugins</groupId>
1010
<artifactId>tis-flink-cdc-kingbase-shade-4-debezium-connector-postgresql</artifactId>
11-
<version>5.0.0</version>
11+
<version>5.1.0</version>
1212
<licenses>
1313
<license>
1414
<name>GNU Affero General Public License</name>
@@ -69,7 +69,7 @@
6969
<dependency>
7070
<groupId>com.qlangtech.tis</groupId>
7171
<artifactId>tis-plugin</artifactId>
72-
<version>5.0.0</version>
72+
<version>5.1.0</version>
7373
<scope>provided</scope>
7474
</dependency>
7575
<dependency>

tis-incr/tis-flink-cdc-mysql-shade-4-debezium-connector-mysql/dependency-reduced-pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
<parent>
44
<artifactId>tis-incr</artifactId>
55
<groupId>com.qlangtech.tis.plugins</groupId>
6-
<version>5.0.0</version>
6+
<version>5.1.0</version>
77
</parent>
88
<modelVersion>4.0.0</modelVersion>
99
<groupId>com.qlangtech.tis.plugins</groupId>
1010
<artifactId>tis-flink-cdc-shade-4-debezium-connector-mysql</artifactId>
11-
<version>5.0.0</version>
11+
<version>5.1.0</version>
1212
<licenses>
1313
<license>
1414
<name>GNU Affero General Public License</name>
@@ -66,7 +66,7 @@
6666
<dependency>
6767
<groupId>com.qlangtech.tis</groupId>
6868
<artifactId>tis-plugin</artifactId>
69-
<version>5.0.0</version>
69+
<version>5.1.0</version>
7070
<scope>provided</scope>
7171
</dependency>
7272
<dependency>

tis-incr/tis-flink-cdc-postgresql-shade-4-debezium-connector-postgresql/dependency-reduced-pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
<parent>
44
<artifactId>tis-incr</artifactId>
55
<groupId>com.qlangtech.tis.plugins</groupId>
6-
<version>5.0.0</version>
6+
<version>5.1.0</version>
77
</parent>
88
<modelVersion>4.0.0</modelVersion>
99
<groupId>com.qlangtech.tis.plugins</groupId>
1010
<artifactId>tis-flink-cdc-postgresql-shade-4-debezium-connector-postgresql</artifactId>
11-
<version>5.0.0</version>
11+
<version>5.1.0</version>
1212
<licenses>
1313
<license>
1414
<name>GNU Affero General Public License</name>
@@ -68,7 +68,7 @@
6868
<dependency>
6969
<groupId>com.qlangtech.tis</groupId>
7070
<artifactId>tis-plugin</artifactId>
71-
<version>5.0.0</version>
71+
<version>5.1.0</version>
7272
<scope>provided</scope>
7373
</dependency>
7474
<dependency>

tis-incr/tis-sink-elasticsearch7/dependency-reduced-pom.xml

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,12 +3,12 @@
33
<parent>
44
<artifactId>tis-incr</artifactId>
55
<groupId>com.qlangtech.tis.plugins</groupId>
6-
<version>5.0.0</version>
6+
<version>5.1.0</version>
77
</parent>
88
<modelVersion>4.0.0</modelVersion>
99
<groupId>com.qlangtech.tis.plugins</groupId>
1010
<artifactId>tis-sink-elasticsearch7</artifactId>
11-
<version>5.0.0</version>
11+
<version>5.1.0</version>
1212
<licenses>
1313
<license>
1414
<name>GNU Affero General Public License</name>
@@ -63,7 +63,7 @@
6363
<dependency>
6464
<groupId>com.qlangtech.tis</groupId>
6565
<artifactId>tis-plugin</artifactId>
66-
<version>5.0.0</version>
66+
<version>5.1.0</version>
6767
<scope>provided</scope>
6868
</dependency>
6969
<dependency>

tis-k8s-plugin/pom.xml

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -68,11 +68,11 @@
6868
<!-- <groupId>tech.powerjob</groupId>-->
6969
<!-- </dependency>-->
7070

71-
<dependency>
72-
<groupId>com.qlangtech.tis.plugins</groupId>
73-
<artifactId>tis-powerjob-common-plugin</artifactId>
74-
<version>${project.version}</version>
75-
</dependency>
71+
<!-- <dependency>-->
72+
<!-- <groupId>com.qlangtech.tis.plugins</groupId>-->
73+
<!-- <artifactId>tis-powerjob-common-plugin</artifactId>-->
74+
<!-- <version>${project.version}</version>-->
75+
<!-- </dependency>-->
7676

7777
<!-- <dependency>-->
7878
<!-- <groupId>io.kubernetes</groupId>-->

tis-k8s-plugin/src/main/java/com/qlangtech/tis/plugin/datax/powerjob/K8SDataXJobWorker.java

Lines changed: 42 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -33,19 +33,15 @@
3333
import com.qlangtech.tis.datax.TimeFormat;
3434
import com.qlangtech.tis.datax.job.DataXJobWorker;
3535
import com.qlangtech.tis.datax.job.ILaunchingOrchestrate;
36-
import com.qlangtech.tis.datax.job.ITISPowerJob;
3736
import com.qlangtech.tis.datax.job.JobOrchestrateException;
3837
import com.qlangtech.tis.datax.job.JobResName;
3938
import com.qlangtech.tis.datax.job.JobResName.OwnerJobExec;
4039
import com.qlangtech.tis.datax.job.PowerjobOrchestrateException;
4140
import com.qlangtech.tis.datax.job.SSERunnable;
4241
import com.qlangtech.tis.datax.job.ServiceResName;
4342
import com.qlangtech.tis.extension.TISExtension;
44-
import com.qlangtech.tis.fullbuild.IFullBuildContext;
4543
import com.qlangtech.tis.fullbuild.indexbuild.RunningStatus;
46-
import com.qlangtech.tis.lang.ErrorValue;
4744
import com.qlangtech.tis.lang.TisException;
48-
import com.qlangtech.tis.lang.TisException.ErrorCode;
4945
import com.qlangtech.tis.plugin.IEndTypeGetter;
5046
import com.qlangtech.tis.plugin.annotation.FormField;
5147
import com.qlangtech.tis.plugin.annotation.FormFieldType;
@@ -106,7 +102,7 @@
106102
* @create: 2021-04-23 18:16
107103
**/
108104
@Public
109-
public class K8SDataXJobWorker extends DataXJobWorker implements ITISPowerJob, ILaunchingOrchestrate {
105+
public class K8SDataXJobWorker extends DataXJobWorker implements ILaunchingOrchestrate {
110106

111107
private static final Logger logger = LoggerFactory.getLogger(K8SDataXJobWorker.class);
112108

@@ -235,10 +231,10 @@ private static ServiceResName<K8SDataXJobWorker> getPowerJobServerService() {
235231

236232
public static final JobResName[] powerJobRes //
237233
= new JobResName[]{
238-
// K8S_DATAX_POWERJOB_MYSQL,
239-
K8S_DATAX_POWERJOB_SERVER
240-
// , K8S_DATAX_POWERJOB_REGISTER_ACCOUNT
241-
//, K8S_DATAX_POWERJOB_WORKER
234+
// K8S_DATAX_POWERJOB_MYSQL,
235+
K8S_DATAX_POWERJOB_SERVER
236+
// , K8S_DATAX_POWERJOB_REGISTER_ACCOUNT
237+
//, K8S_DATAX_POWERJOB_WORKER
242238
};
243239

244240

@@ -263,17 +259,17 @@ private static ServiceResName<K8SDataXJobWorker> getPowerJobServerService() {
263259
private transient CoreV1Api apiClient;
264260

265261

266-
public String getPowerJobMasterGateway() {
267-
try {
268-
final String linkHost = this.serverPortExport
269-
.getClusterHost(this.getK8SApi(), this.getImage(), powerJobServiceResAndOwnerGetter.get());
270-
return linkHost;
271-
} catch (ServiceNotDefinedException e) {
272-
//
273-
throw throwPowerJobClusterLossOfContactException(Optional.of(e));
274-
// throw new RuntimeException(e);
275-
}
276-
}
262+
// public String getPowerJobMasterGateway() {
263+
// try {
264+
// final String linkHost = this.serverPortExport
265+
// .getClusterHost(this.getK8SApi(), this.getImage(), powerJobServiceResAndOwnerGetter.get());
266+
// return linkHost;
267+
// } catch (ServiceNotDefinedException e) {
268+
// //
269+
// throw throwPowerJobClusterLossOfContactException(Optional.of(e));
270+
// // throw new RuntimeException(e);
271+
// }
272+
// }
277273

278274
public CoreV1Api getK8SApi() {
279275
if (this.apiClient == null) {
@@ -289,21 +285,21 @@ public CoreV1Api getK8SApi() {
289285
// public static String getDefaultZookeeperAddress() {
290286
// return processDefaultHost(Config.getZKHost());
291287
// }
292-
private transient TISPowerJobClient powerJobClient;
293-
294-
@Override
295-
public TISPowerJobClient getPowerJobClient() {
296-
if (powerJobClient == null) {
297-
try {
298-
powerJobClient = TISPowerJobClient.create(
299-
this.serverPortExport.getClusterHost(this.getK8SApi(), this.getK8SImage(), powerJobServiceResAndOwnerGetter.get())
300-
, this.appName, this.password);
301-
} catch (ServiceNotDefinedException e) {
302-
throw throwPowerJobClusterLossOfContactException(Optional.of(e));
303-
}
304-
}
305-
return powerJobClient;
306-
}
288+
// private transient TISPowerJobClient powerJobClient;
289+
//
290+
// @Override
291+
// public TISPowerJobClient getPowerJobClient() {
292+
// if (powerJobClient == null) {
293+
// try {
294+
// powerJobClient = TISPowerJobClient.create(
295+
// this.serverPortExport.getClusterHost(this.getK8SApi(), this.getK8SImage(), powerJobServiceResAndOwnerGetter.get())
296+
// , this.appName, this.password);
297+
// } catch (ServiceNotDefinedException e) {
298+
// throw throwPowerJobClusterLossOfContactException(Optional.of(e));
299+
// }
300+
// }
301+
// return powerJobClient;
302+
// }
307303

308304
@Override
309305
public Map<String, Object> getPayloadInfo() {
@@ -316,8 +312,8 @@ public Map<String, Object> getPayloadInfo() {
316312
// throwPowerJobClusterLossOfContactException();
317313
return payloads;
318314
} catch (ServiceNotDefinedException e) {
319-
// throw new RuntimeException(e);
320-
throw throwPowerJobClusterLossOfContactException(Optional.of(e));
315+
throw new RuntimeException(e);
316+
// throw throwPowerJobClusterLossOfContactException(Optional.of(e));
321317
}
322318
}
323319

@@ -590,8 +586,8 @@ public List<RcDeployment> getRCDeployments() {
590586
K8SController k8SController = getK8SController();
591587
RcDeployment powerjobServer = k8SController.getRCDeployment(K8S_DATAX_POWERJOB_SERVER);
592588
if (powerjobServer == null) {
593-
// throw TisException.create("the powerJob has been loss of communication");
594-
throw throwPowerJobClusterLossOfContactException(Optional.empty());
589+
throw TisException.create("the powerJob has been loss of communication");
590+
// throw throwPowerJobClusterLossOfContactException(Optional.empty());
595591
}
596592
powerjobServer.setReplicaScalable(false);
597593

@@ -612,12 +608,12 @@ public List<RcDeployment> getRCDeployments() {
612608
// return getK8SController().getRCDeployment(DataXJobWorker.K8S_DATAX_INSTANCE_NAME);
613609
}
614610

615-
private static TisException throwPowerJobClusterLossOfContactException(Optional<ServiceNotDefinedException> e) {
616-
return TisException.create(
617-
ErrorValue.create(ErrorCode.POWER_JOB_CLUSTER_LOSS_OF_CONTACT
618-
, IFullBuildContext.KEY_TARGET_NAME, TargetResName.K8S_DATAX_INSTANCE_NAME.getName())
619-
, e.map((except) -> except.getMessage()).orElse("the powerJob has been loss of communication"));
620-
}
611+
// private static TisException throwPowerJobClusterLossOfContactException(Optional<ServiceNotDefinedException> e) {
612+
// return TisException.create(
613+
// ErrorValue.create(ErrorCode.POWER_JOB_CLUSTER_LOSS_OF_CONTACT
614+
// , IFullBuildContext.KEY_TARGET_NAME, TargetResName.K8S_DATAX_INSTANCE_NAME.getName())
615+
// , e.map((except) -> except.getMessage()).orElse("the powerJob has been loss of communication"));
616+
// }
621617

622618
@Override
623619
public WatchPodLog listPodAndWatchLog(String podName, ILogListener listener) {
@@ -750,8 +746,6 @@ public final PowerJobK8SImage getPowerJobImage() {
750746
// protected K8SDataXPowerJobWorker getPowerJobWorker() {
751747
// return K8SUtils.getK8SDataXPowerJobWorker();
752748
// }
753-
754-
755749
public NamespacedEventCallCriteria launchPowerjobServer() throws ApiException, PowerjobOrchestrateException {
756750
SSERunnable sse = SSERunnable.getLocal();
757751

@@ -804,7 +798,7 @@ public NamespacedEventCallCriteria launchPowerjobServer() throws ApiException, P
804798
// }
805799
// --spring.datasource.core.jdbc-url=
806800
// --oms.transporter.active.protocols=http
807-
// envVar.setValue(" --oms.mongodb.enable=false " + coreJbdcParams);
801+
// envVar.setValue(" --oms.mongodb.enable=false " + coreJbdcParams);
808802
envs.add(envVar);
809803

810804

0 commit comments

Comments
 (0)