Skip to content

Commit 4a564f5

Browse files
committed
add a processing for flinkCluster available telnet
1 parent a631d0f commit 4a564f5

10 files changed

Lines changed: 103 additions & 22 deletions

File tree

tis-datax/executor/dolphinscheduler-task-tis-datasync/dependency-reduced-pom.xml

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,16 @@
7070
<artifactId>tis-base-test</artifactId>
7171
<version>4.3.0-SNAPSHOT</version>
7272
<scope>test</scope>
73+
<exclusions>
74+
<exclusion>
75+
<artifactId>mockito-core</artifactId>
76+
<groupId>org.mockito</groupId>
77+
</exclusion>
78+
<exclusion>
79+
<artifactId>mockito-inline</artifactId>
80+
<groupId>org.mockito</groupId>
81+
</exclusion>
82+
</exclusions>
7383
</dependency>
7484
<dependency>
7585
<groupId>org.apache.dolphinscheduler</groupId>

tis-incr/tis-flink-cdc-common-plugin-shade-4-debezium/pom.xml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,9 @@
6464
<exclude>log4j:log4j</exclude>
6565
<exclude>org.slf4j:slf4j-api</exclude>
6666
<exclude>org.codehaus.groovy:groovy-all</exclude>
67+
<exclude>org.apache.flink:flink-cdc-common</exclude>
68+
<exclude>org.apache.flink:flink-cdc-runtime</exclude>
69+
<exclude>org.apache.flink:flink-shaded-guava</exclude>
6770
</excludes>
6871
</artifactSet>
6972
<filters>

tis-incr/tis-flink-cdc-common/pom.xml

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,36 @@
5656
</exclusions>
5757
</dependency>
5858

59+
<!-- <dependency>-->
60+
<!-- <groupId>org.apache.flink</groupId>-->
61+
<!-- <artifactId>flink-cdc-common</artifactId>-->
62+
<!-- <version>${flink.cdc.version}</version>-->
63+
<!-- </dependency>-->
64+
65+
<!-- <dependency>-->
66+
<!-- <groupId>org.apache.flink</groupId>-->
67+
<!-- <artifactId>flink-cdc-runtime</artifactId>-->
68+
<!-- <version>${flink.cdc.version}</version>-->
69+
<!-- <exclusions>-->
70+
<!-- <exclusion>-->
71+
<!-- <groupId>org.apache.calcite</groupId>-->
72+
<!-- <artifactId>calcite-core</artifactId>-->
73+
<!-- </exclusion>-->
74+
<!-- </exclusions>-->
75+
<!-- </dependency>-->
76+
77+
<!-- <dependency>-->
78+
<!-- <groupId>org.apache.flink</groupId>-->
79+
<!-- <artifactId>flink-cdc-composer-tis</artifactId>-->
80+
<!-- <version>${flink.cdc.version}</version>-->
81+
<!-- <exclusions>-->
82+
<!-- <exclusion>-->
83+
<!-- <groupId>org.apache.flink</groupId>-->
84+
<!-- <artifactId>flink-kubernetes</artifactId>-->
85+
<!-- </exclusion>-->
86+
<!-- </exclusions>-->
87+
<!-- </dependency>-->
88+
5989
<!-- <dependency>-->
6090
<!-- <groupId>org.apache.flink</groupId>-->
6191
<!-- <artifactId>flink-connector-debezium</artifactId>-->

tis-incr/tis-flink-cdc-mysql-plugin/pom.xml

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,14 @@
6969
<groupId>org.apache.flink</groupId>
7070
<artifactId>flink-connector-debezium</artifactId>
7171
</exclusion>
72+
<exclusion>
73+
<groupId>org.apache.flink</groupId>
74+
<artifactId>flink-cdc-common</artifactId>
75+
</exclusion>
76+
<exclusion>
77+
<groupId>org.apache.flink</groupId>
78+
<artifactId>flink-cdc-runtime</artifactId>
79+
</exclusion>
7280
</exclusions>
7381
</dependency>
7482

tis-incr/tis-flink-cdc-oracle-plugin/pom.xml

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,14 @@
7171
<groupId>org.apache.flink</groupId>
7272
<artifactId>flink-connector-debezium</artifactId>
7373
</exclusion>
74+
<exclusion>
75+
<groupId>org.apache.flink</groupId>
76+
<artifactId>flink-cdc-common</artifactId>
77+
</exclusion>
78+
<exclusion>
79+
<groupId>org.apache.flink</groupId>
80+
<artifactId>flink-cdc-runtime</artifactId>
81+
</exclusion>
7482
<exclusion>
7583
<groupId>io.debezium</groupId>
7684
<artifactId>debezium-core</artifactId>

tis-incr/tis-flink-cdc-postgresql-plugin/pom.xml

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,14 @@
7171
<groupId>org.apache.flink</groupId>
7272
<artifactId>flink-connector-debezium</artifactId>
7373
</exclusion>
74+
<exclusion>
75+
<groupId>org.apache.flink</groupId>
76+
<artifactId>flink-cdc-common</artifactId>
77+
</exclusion>
78+
<exclusion>
79+
<groupId>org.apache.flink</groupId>
80+
<artifactId>flink-cdc-runtime</artifactId>
81+
</exclusion>
7482
<exclusion>
7583
<groupId>io.debezium</groupId>
7684
<artifactId>debezium-connector-postgres</artifactId>

tis-incr/tis-realtime-flink/src/main/java/com/qlangtech/plugins/incr/flink/common/FlinkCluster.java

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333
import com.qlangtech.tis.plugin.annotation.FormField;
3434
import com.qlangtech.tis.plugin.annotation.FormFieldType;
3535
import com.qlangtech.tis.plugin.annotation.Validator;
36+
import com.qlangtech.tis.realtime.utils.NetUtils;
3637
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
3738
import org.apache.flink.client.program.ClusterClient;
3839
import org.apache.flink.client.program.rest.RestClusterClient;
@@ -43,6 +44,7 @@
4344
import org.slf4j.Logger;
4445
import org.slf4j.LoggerFactory;
4546

47+
import java.net.SocketTimeoutException;
4648
import java.time.Duration;
4749
import java.util.Collection;
4850
import java.util.Optional;
@@ -95,14 +97,27 @@ public ClusterClient createConfigInstance() {
9597
return createFlinkRestClusterClient(Optional.empty(), Optional.of(2000l));
9698
}
9799

100+
public static Configuration setNoAttamptsForClient(Configuration configuration, Long timeoutMillis) {
101+
configuration.set(RestOptions.CONNECTION_TIMEOUT, Duration.ofMillis(timeoutMillis));
102+
configuration.setInteger(RestOptions.RETRY_MAX_ATTEMPTS, 0);
103+
configuration.set(RestOptions.RETRY_DELAY, Duration.ofMillis(0l));
104+
return configuration;
105+
}
106+
98107
/**
99108
* @param connTimeout The maximum time in ms for the client to establish a TCP connection.
100109
* @return
101110
*/
102111
public ClusterClient createFlinkRestClusterClient(Optional<String> clusterId, Optional<Long> connTimeout) {
112+
JobManagerAddress managerAddress = this.getJobManagerAddress();
113+
try {
114+
managerAddress.telnet();
115+
} catch (SocketTimeoutException e) {
116+
throw new RuntimeException(e);
117+
}
103118

104119
try {
105-
JobManagerAddress managerAddress = this.getJobManagerAddress();
120+
106121
Configuration configuration = new Configuration();
107122
configuration.setString(JobManagerOptions.ADDRESS, managerAddress.host);
108123
configuration.setInteger(JobManagerOptions.PORT, managerAddress.port);
@@ -111,9 +126,10 @@ public ClusterClient createFlinkRestClusterClient(Optional<String> clusterId, Op
111126
configuration.set(RestOptions.RETRY_DELAY, Duration.ofMillis(this.retryDelay));
112127

113128
if (connTimeout.isPresent()) {
114-
configuration.set(RestOptions.CONNECTION_TIMEOUT, Duration.ofSeconds(connTimeout.get()));
115-
configuration.setInteger(RestOptions.RETRY_MAX_ATTEMPTS, 0);
116-
configuration.set(RestOptions.RETRY_DELAY, Duration.ofMillis(0l));
129+
setNoAttamptsForClient(configuration, connTimeout.get());
130+
// configuration.set(RestOptions.CONNECTION_TIMEOUT, Duration.ofSeconds(connTimeout.get()));
131+
// configuration.setInteger(RestOptions.RETRY_MAX_ATTEMPTS, 0);
132+
// configuration.set(RestOptions.RETRY_DELAY, Duration.ofMillis(0l));
117133
}
118134

119135

tis-incr/tis-realtime-flink/src/main/java/com/qlangtech/plugins/incr/flink/launch/StateBackendFactory.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ public static Optional<IncrStreamFactory.ISavePointSupport> getSavePointSupport(
4949

5050

5151
/**
52-
* 缺的当前执行任务的状态
52+
* 取得当前执行任务的状态
5353
*
5454
* @return
5555
*/

tis-k8s-plugin/src/main/java/com/qlangtech/tis/config/k8s/impl/DefaultK8sContext.java

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import com.alibaba.citrus.turbine.Context;
2121
import com.qlangtech.tis.annotation.Public;
2222
import com.qlangtech.tis.config.ParamsConfig;
23+
import com.qlangtech.tis.config.flink.JobManagerAddress;
2324
import com.qlangtech.tis.config.k8s.IK8sContext;
2425
import com.qlangtech.tis.extension.Descriptor;
2526
import com.qlangtech.tis.extension.TISExtension;
@@ -43,6 +44,9 @@
4344

4445
import java.io.Reader;
4546
import java.io.StringReader;
47+
import java.net.MalformedURLException;
48+
import java.net.SocketTimeoutException;
49+
import java.net.URL;
4650
import java.util.Map;
4751
import java.util.stream.Collectors;
4852

@@ -88,6 +92,15 @@ public String getKubeBasePath() {
8892
public ApiClient createConfigInstance() {
8993

9094
ApiClient client = null;
95+
URL basePath = null;
96+
try {
97+
basePath = new URL(this.kubeBasePath);
98+
(new JobManagerAddress(basePath.getHost(), basePath.getPort())).telnet();
99+
} catch (MalformedURLException e) {
100+
throw new RuntimeException(e);
101+
} catch (SocketTimeoutException e) {
102+
throw new RuntimeException("basePath:" + basePath, e);
103+
}
91104
try {
92105
try (Reader reader = new StringReader(this.kubeConfigContent)) {
93106
client = ClientBuilder.kubeconfig(KubeConfig.loadKubeConfig(reader)).setBasePath(this.kubeBasePath).build();
@@ -138,8 +151,8 @@ protected boolean verify(IControlMsgHandler msgHandler, Context context, PostFor
138151
} else {
139152
msgHandler.addActionMessage(context
140153
, "exist namespace is:" + namespaceList.getItems().stream().map((ns) -> {
141-
return ns.getMetadata().getName();
142-
}).collect(Collectors.joining(",")));
154+
return ns.getMetadata().getName();
155+
}).collect(Collectors.joining(",")));
143156
}
144157
} catch (Throwable e) {
145158
logger.warn(e.getMessage(), e);

tis-k8s-plugin/src/main/java/com/qlangtech/tis/plugin/k8s/K8sExceptionUtils.java

Lines changed: 0 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -40,21 +40,6 @@ public static TisException convert(String msg, ApiException e) {
4040
}
4141

4242
public static TisException convert(ErrorValue errCode, String msg, ApiException e) {
43-
44-
// final ClassLoader current = Thread.currentThread().getContextClassLoader();
45-
// try {
46-
// Thread.currentThread().setContextClassLoader(V1Status.class.getClassLoader());
47-
// V1Status v1Status = JSON.parseObject(e.getResponseBody(), V1Status.class);
48-
// String errMsg = msg;
49-
// if (v1Status != null) {
50-
// errMsg = (msg == null) ? v1Status.getMessage() : msg + ":" + v1Status.getMessage();
51-
// }
52-
// return TisException.create(errCode, StringUtils.defaultIfEmpty(errMsg, e.getMessage()), e);
53-
//
54-
// } finally {
55-
// Thread.currentThread().setContextClassLoader(current);
56-
// }
57-
5843
try {
5944
return ClassloaderUtils.processByResetThreadClassloader(V1Status.class, () -> {
6045
V1Status v1Status = JSON.parseObject(e.getResponseBody(), V1Status.class);

0 commit comments

Comments
 (0)