Skip to content

Commit e27f480

Browse files
committed
add paimon batch and incr sink support for TIS
1 parent db2dda2 commit e27f480

11 files changed

Lines changed: 174 additions & 64 deletions

File tree

tis-datax/tis-datax-hdfs-plugin/src/main/java/com/qlangtech/tis/hdfs/impl/HdfsFileSystemFactory.java

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -229,11 +229,11 @@ private static Configuration getConfiguration(String hdfsSiteContent, Boolean us
229229
// 这个缓存还是需要的,不然如果另外的调用FileSystem实例不是通过调用getFileSystem这个方法的进入,就调用不到了
230230
conf.setBoolean("fs.hdfs.impl.disable.cache", false);
231231

232-
// if (StringUtils.isNotEmpty(this.kerberos)) {
233-
// Logger.info("kerberos has been enabled,name:" + this.kerberos);
234-
// KerberosCfg kerberosCfg = KerberosCfg.getKerberosCfg(this.kerberos);
235-
// kerberosCfg.setConfiguration(conf);
236-
// }
232+
// if (StringUtils.isNotEmpty(this.kerberos)) {
233+
// Logger.info("kerberos has been enabled,name:" + this.kerberos);
234+
// KerberosCfg kerberosCfg = KerberosCfg.getKerberosCfg(this.kerberos);
235+
// kerberosCfg.setConfiguration(conf);
236+
// }
237237
cfgProcess.accept(conf);
238238
// this.setConfiguration(conf);
239239
conf.reloadConfiguration();
@@ -285,6 +285,14 @@ public static <T> T setConfiguration(
285285
}
286286
}
287287

288+
@Override
289+
public String getRootDir() {
290+
if (StringUtils.isEmpty(this.rootDir)) {
291+
throw new IllegalStateException("prop rootDir can not be null");
292+
}
293+
return rootDir;
294+
}
295+
288296
@Override
289297
public String getFSAddress() {
290298
return getConfiguration().get(CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY);

tis-datax/tis-datax-hdfs-plugin/src/main/java/com/qlangtech/tis/plugin/datax/BasicFSWriter.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,7 +119,7 @@ public final IDataxContext getSubTask(Optional<IDataxProcessor.TableMap> tableMa
119119

120120
protected abstract FSDataXContext getDataXContext(IDataxProcessor.TableMap tableMap, Optional<RecordTransformerRules> transformerRules);
121121

122-
public class FSDataXContext implements IDataxContext {
122+
public class FSDataXContext implements IDataxContext {
123123

124124
protected final IDataxProcessor.TableMap tabMap;
125125
private final String dataxName;

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

Lines changed: 30 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
package com.qlangtech.tis.plugin.common;
2020

2121
import com.qlangtech.tis.extension.Describable;
22+
import com.qlangtech.tis.extension.Descriptor;
2223
import com.qlangtech.tis.trigger.util.JsonUtil;
2324
import com.qlangtech.tis.util.DescriptorsJSON;
2425
import org.junit.Assert;
@@ -31,16 +32,38 @@ public class PluginDesc {
3132

3233

3334
public static <TT extends Describable> void testDescGenerate(Class<TT> clazz, String assertFileName) {
35+
// try {
3436
try {
3537
TT plugin = clazz.newInstance();
36-
// DescriptorsJSON descJson = new DescriptorsJSON(plugin.getDescriptor());
37-
JsonUtil.assertJSONEqual(clazz, assertFileName
38-
, JsonUtil.toString( DescriptorsJSON.desc(plugin.getDescriptor())), (m, e, a) -> {
39-
Assert.assertEquals(m, e, a);
40-
});
41-
//return plugin;
38+
Descriptor descriptor = plugin.getDescriptor();
39+
if (descriptor == null) {
40+
throw new NullPointerException("plugin:" + plugin.getClass().getName() + " relevant descriptor can not be null");
41+
}
42+
43+
testDescGenerate(clazz, descriptor, assertFileName);
4244
} catch (Exception e) {
43-
throw new RuntimeException(assertFileName, e);
45+
throw new RuntimeException(e);
46+
}
47+
//return plugin;
48+
// } catch (Exception e) {
49+
// throw new RuntimeException(assertFileName, e);
50+
// }
51+
}
52+
53+
public static <TT extends Describable> void testDescGenerate(Class<TT> clazz, Descriptor descriptor, String assertFileName) {
54+
// try {
55+
56+
if (descriptor == null) {
57+
throw new NullPointerException("plugin:" + clazz.getName() + " relevant descriptor can not be null");
4458
}
59+
// DescriptorsJSON descJson = new DescriptorsJSON(plugin.getDescriptor());
60+
JsonUtil.assertJSONEqual(clazz, assertFileName
61+
, JsonUtil.toString(DescriptorsJSON.desc(descriptor)), (m, e, a) -> {
62+
Assert.assertEquals(m, e, a);
63+
});
64+
//return plugin;
65+
// } catch (Exception e) {
66+
// throw new RuntimeException(assertFileName, e);
67+
// }
4568
}
4669
}

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/plugins/incr/flink/TISChildFirstClassLoader.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,9 @@
1919
package com.qlangtech.plugins.incr.flink;
2020

2121
import com.qlangtech.tis.extension.UberClassLoader;
22+
import org.apache.flink.util.ChildFirstClassLoader;
2223
import org.apache.flink.util.FlinkUserCodeClassLoader;
24+
import org.apache.flink.util.MutableURLClassLoader;
2325

2426
import java.io.IOException;
2527
import java.net.URL;
@@ -132,4 +134,9 @@ public URL nextElement() {
132134
ClassLoader.registerAsParallelCapable();
133135
}
134136

137+
@Override
138+
public MutableURLClassLoader copy() {
139+
return new TISChildFirstClassLoader(
140+
this.uberClassloader, this.getURLs(), this.getParent() ,alwaysParentFirstPatterns, classLoadingExceptionHandler);
141+
}
135142
}

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/plugins/incr/flink/cdc/FlinkCDCPipelineEventProcess.java

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,14 @@
1818

1919
package com.qlangtech.plugins.incr.flink.cdc;
2020

21+
2122
import java.math.BigDecimal;
2223
import java.time.LocalDateTime;
2324

2425
/**
2526
* @author: 百岁(baisui@qlangtech.com)
2627
* @create: 2025-05-23 10:37
28+
* @see org.apache.flink.cdc.runtime.serializer.data.writer.BinaryWriter#write
2729
**/
2830
public class FlinkCDCPipelineEventProcess extends BiFunction {
2931

@@ -44,6 +46,22 @@ public org.apache.flink.cdc.common.types.DataType getDataType() {
4446
return this.dataType;
4547
}
4648

49+
50+
public static class FlinkPipelineStringConvert extends BiFunction {
51+
@Override
52+
public Object deApply(Object o) {
53+
throw new UnsupportedOperationException();
54+
}
55+
56+
@Override
57+
public Object apply(Object o) {
58+
if (o instanceof java.nio.ByteBuffer) {
59+
return org.apache.flink.cdc.common.data.binary.BinaryStringData.fromBytes(((java.nio.ByteBuffer) o).array());
60+
}
61+
return org.apache.flink.cdc.common.data.binary.BinaryStringData.fromString(String.valueOf(o));
62+
}
63+
}
64+
4765
public static class FlinkPipelineDecimalConvert extends BiFunction {
4866
// private final DataType type;
4967

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/tis/realtime/DTOSourceTagProcessFunction.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -31,8 +31,8 @@
3131
public class DTOSourceTagProcessFunction extends SourceProcessFunction<DTO> {
3232
public static final String KEY_MERGE_ALL_TABS_IN_ONE_BUS = "merge_all_tabs_in_one_bus";
3333

34-
public static DTOSourceTagProcessFunction createMergeAllTabsInOneBus() {
35-
return new MergeAllTabsInOneBusProcessFunction();
34+
public static DTOSourceTagProcessFunction create(boolean flinkCDCPipelineEnable, Map<String, OutputTag<DTO>> tab2OutputTag) {
35+
return flinkCDCPipelineEnable ? new MergeAllTabsInOneBusProcessFunction() : new DTOSourceTagProcessFunction(tab2OutputTag);
3636
}
3737

3838
public DTOSourceTagProcessFunction(Map<String, OutputTag<DTO>> tab2OutputTag) {
@@ -47,7 +47,7 @@ protected String getTableName(DTO record) {
4747

4848
static class MergeAllTabsInOneBusProcessFunction extends DTOSourceTagProcessFunction {
4949

50-
public MergeAllTabsInOneBusProcessFunction() {
50+
private MergeAllTabsInOneBusProcessFunction() {
5151
super(Collections.singletonMap(KEY_MERGE_ALL_TABS_IN_ONE_BUS, new OutputTag<DTO>(KEY_MERGE_ALL_TABS_IN_ONE_BUS) {
5252
}));
5353
}

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/tis/realtime/ReaderSource.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -91,12 +91,12 @@ protected DataStreamSource<DTO> addAsSource(StreamExecutionEnvironment env) {
9191

9292
@Override
9393
protected SourceProcessFunction<DTO> createStreamTagFunction(Map<String, OutputTag<DTO>> tab2OutputTag) {
94-
return flinkCDCPipelineEnable ? DTOSourceTagProcessFunction.createMergeAllTabsInOneBus() : new DTOSourceTagProcessFunction(tab2OutputTag);
94+
return DTOSourceTagProcessFunction.create(flinkCDCPipelineEnable, tab2OutputTag);// : new DTOSourceTagProcessFunction(tab2OutputTag);
9595
}
9696
};
9797
}
9898

99-
public static ReaderSource<DTO> createDTOSource(String tokenName, final DataStreamSource<DTO> source) {
99+
public static ReaderSource<DTO> createDTOSource(String tokenName, boolean flinkCDCPipelineEnable, final DataStreamSource<DTO> source) {
100100
return new SideOutputReaderSource<DTO>(tokenName) {
101101
@Override
102102
protected DataStreamSource<DTO> addAsSource(StreamExecutionEnvironment env) {
@@ -105,7 +105,7 @@ protected DataStreamSource<DTO> addAsSource(StreamExecutionEnvironment env) {
105105

106106
@Override
107107
protected SourceProcessFunction<DTO> createStreamTagFunction(Map<String, OutputTag<DTO>> tab2OutputTag) {
108-
return new DTOSourceTagProcessFunction(tab2OutputTag);
108+
return DTOSourceTagProcessFunction.create(flinkCDCPipelineEnable, tab2OutputTag);// new DTOSourceTagProcessFunction(tab2OutputTag);
109109
}
110110
};
111111
}

tis-incr/tis-flink-extends/src/main/java/org/apache/flink/kubernetes/entrypoint/KubernetesApplicationClusterEntrypointOfTIS.java

Lines changed: 73 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -23,29 +23,40 @@
2323
import com.qlangtech.tis.coredefine.module.action.TargetResName;
2424
import com.qlangtech.tis.manage.common.incr.StreamContextConstant.TISRes;
2525
import com.qlangtech.tis.manage.common.incr.UberJarUtil;
26+
import org.apache.flink.client.cli.ArtifactFetchOptions;
2627
import org.apache.flink.client.deployment.application.ApplicationClusterEntryPoint;
2728
import org.apache.flink.client.deployment.application.ApplicationConfiguration;
28-
import org.apache.flink.client.deployment.application.ClassPathPackagedProgramRetriever;
29+
import org.apache.flink.client.program.DefaultPackagedProgramRetriever;
2930
import org.apache.flink.client.program.PackagedProgram;
3031
import org.apache.flink.client.program.PackagedProgramRetriever;
3132
import org.apache.flink.client.program.PackagedProgramUtils;
33+
import org.apache.flink.client.program.artifact.ArtifactFetchManager;
3234
import org.apache.flink.configuration.Configuration;
3335
import org.apache.flink.configuration.PipelineOptions;
34-
import org.apache.flink.kubernetes.utils.KubernetesUtils;
36+
import org.apache.flink.core.fs.FileSystem;
37+
import org.apache.flink.core.plugin.PluginManager;
38+
import org.apache.flink.core.plugin.PluginUtils;
39+
import org.apache.flink.kubernetes.configuration.KubernetesConfigOptions;
3540
import org.apache.flink.runtime.entrypoint.ClusterEntrypoint;
3641
import org.apache.flink.runtime.entrypoint.ClusterEntrypointUtils;
3742
import org.apache.flink.runtime.entrypoint.DynamicParametersConfigurationParserFactory;
43+
import org.apache.flink.runtime.security.contexts.SecurityContext;
44+
import org.apache.flink.runtime.util.EnvironmentInformation;
45+
import org.apache.flink.runtime.util.JvmShutdownSafeguard;
46+
import org.apache.flink.runtime.util.SignalHandler;
3847
import org.apache.flink.util.FlinkException;
39-
import org.apache.flink.util.Preconditions;
4048
import org.slf4j.Logger;
4149
import org.slf4j.LoggerFactory;
4250

51+
import javax.annotation.Nullable;
4352
import java.io.File;
44-
import java.io.IOException;
4553
import java.net.URL;
54+
import java.util.Collections;
4655
import java.util.List;
4756
import java.util.Objects;
4857

58+
import static org.apache.flink.util.Preconditions.checkArgument;
59+
4960
/**
5061
* 当客户端选择使用 kerbernetes-application 部署方式的时候,在flink-jobManager 端服务组装构建jobGraph实例,需要从TIS-console端拉取 同步任务对应的jar包到本地
5162
*
@@ -64,12 +75,16 @@ private KubernetesApplicationClusterEntrypointOfTIS(Configuration configuration,
6475

6576
public static void main(final String[] args) {
6677
// startup checks and logging
78+
EnvironmentInformation.logEnvironmentInfo(
79+
LOG, KubernetesApplicationClusterEntrypoint.class.getSimpleName(), args);
80+
SignalHandler.register(LOG);
81+
JvmShutdownSafeguard.installAsShutdownHook(LOG);
6782

6883
final Configuration dynamicParameters =
6984
ClusterEntrypointUtils.parseParametersOrExit(
7085
args,
7186
new DynamicParametersConfigurationParserFactory(),
72-
KubernetesApplicationClusterEntrypointOfTIS.class);
87+
KubernetesApplicationClusterEntrypoint.class);
7388
final Configuration configuration =
7489
KubernetesEntrypointUtils.loadConfiguration(dynamicParameters);
7590

@@ -80,28 +95,33 @@ public static void main(final String[] args) {
8095
break;
8196
}
8297
Objects.requireNonNull(targetResName, "targetResName can not be null");
98+
8399
PackagedProgram program = null;
84100
try {
85-
// baisui modify 2024/1/8
86-
// 下载最新的jar包
87-
// PluginMeta flinkPluginMeta = TISFlinkClassLoaderFactory.getFlinkPluginMeta(targetResName);
88-
// flinkPluginMeta.copyFromRemote();
101+
89102
TISRes unberJarFile = UberJarUtil.getStreamUberJarFile(targetResName);
90103
unberJarFile.sync2Local(true);
91104
String unberJarURL = "local" + String.valueOf(unberJarFile.getFile().toURI().toURL()).substring("file".length());
92105
logger.info("TIS unberJarURL:{}", unberJarURL);
93106
TISFlinkClassLoaderFactory.synAppRelevantCfgsAndTpis(new URL[]{unberJarFile.getFile().toURI().toURL()});
94107
configuration.set(PipelineOptions.JARS, Lists.newArrayList(unberJarURL));
95-
program = getPackagedProgram(configuration);
108+
109+
PluginManager pluginManager =
110+
PluginUtils.createPluginManagerFromRootFolder(configuration);
111+
LOG.info(
112+
"Install default filesystem for fetching user artifacts in Kubernetes Application Mode.");
113+
FileSystem.initialize(configuration, pluginManager);
114+
SecurityContext securityContext = installSecurityContext(configuration);
115+
program = securityContext.runSecured(() -> getPackagedProgram(configuration));
96116
} catch (Exception e) {
97-
logger.error("Could not create application program.", e);
117+
LOG.error("Could not create application program.", e);
98118
System.exit(1);
99119
}
100120

101121
try {
102122
configureExecution(configuration, program);
103123
} catch (Exception e) {
104-
logger.error("Could not apply application configuration.", e);
124+
LOG.error("Could not apply application configuration.", e);
105125
System.exit(1);
106126
}
107127

@@ -112,7 +132,7 @@ public static void main(final String[] args) {
112132
}
113133

114134
private static PackagedProgram getPackagedProgram(final Configuration configuration)
115-
throws IOException, FlinkException {
135+
throws FlinkException {
116136

117137
final ApplicationConfiguration applicationConfiguration =
118138
ApplicationConfiguration.fromConfiguration(configuration);
@@ -128,24 +148,53 @@ private static PackagedProgram getPackagedProgram(final Configuration configurat
128148
private static PackagedProgramRetriever getPackagedProgramRetriever(
129149
final Configuration configuration,
130150
final String[] programArguments,
131-
final String jobClassName)
132-
throws IOException {
151+
@Nullable final String jobClassName)
152+
throws FlinkException {
133153

134154
final File userLibDir = ClusterEntrypointUtils.tryFindUserLibDirectory().orElse(null);
135-
final ClassPathPackagedProgramRetriever.Builder retrieverBuilder =
136-
ClassPathPackagedProgramRetriever.newBuilder(programArguments)
137-
.setUserLibDirectory(userLibDir)
138-
.setJobClassName(jobClassName);
139155

140156
// No need to do pipelineJars validation if it is a PyFlink job.
141157
if (!(PackagedProgramUtils.isPython(jobClassName)
142158
|| PackagedProgramUtils.isPython(programArguments))) {
143-
final List<File> pipelineJars =
144-
KubernetesUtils.checkJarFileForApplicationMode(configuration);
145-
Preconditions.checkArgument(pipelineJars.size() == 1, "Should only have one jar");
146-
retrieverBuilder.setJarFile(pipelineJars.get(0));
159+
final ArtifactFetchManager.Result fetchRes = fetchArtifacts(configuration);
160+
161+
return DefaultPackagedProgramRetriever.create(
162+
userLibDir,
163+
fetchRes.getJobJar(),
164+
fetchRes.getArtifacts(),
165+
jobClassName,
166+
programArguments,
167+
configuration);
168+
}
169+
170+
return DefaultPackagedProgramRetriever.create(
171+
userLibDir, jobClassName, programArguments, configuration);
172+
}
173+
174+
private static ArtifactFetchManager.Result fetchArtifacts(Configuration configuration) {
175+
try {
176+
String targetDir = generateJarDir(configuration);
177+
ArtifactFetchManager fetchMgr = new ArtifactFetchManager(configuration, targetDir);
178+
179+
List<String> uris = configuration.get(PipelineOptions.JARS);
180+
checkArgument(uris.size() == 1, "Should only have one jar");
181+
List<String> additionalUris =
182+
configuration
183+
.getOptional(ArtifactFetchOptions.ARTIFACT_LIST)
184+
.orElse(Collections.emptyList());
185+
186+
return fetchMgr.fetchArtifacts(uris.get(0), additionalUris);
187+
} catch (Exception e) {
188+
throw new RuntimeException(e);
147189
}
148-
return retrieverBuilder.build();
190+
}
191+
192+
static String generateJarDir(Configuration configuration) {
193+
return String.join(
194+
File.separator,
195+
new File(configuration.get(ArtifactFetchOptions.BASE_DIR)).getAbsolutePath(),
196+
configuration.get(KubernetesConfigOptions.NAMESPACE),
197+
configuration.get(KubernetesConfigOptions.CLUSTER_ID));
149198
}
150199

151200

0 commit comments

Comments
 (0)