Skip to content

Commit ff31e9d

Browse files
committed
add pipeline sink supporting for TIS kafka connector ,issue:datavane/tis#484
1 parent 1569e37 commit ff31e9d

15 files changed

Lines changed: 1172 additions & 13 deletions

File tree

tis-incr/tis-flink-chunjun-kafka-plugin/pom.xml

Lines changed: 23 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,12 @@
7878
</exclusions>
7979
</dependency>
8080

81+
<dependency>
82+
<groupId>org.apache.flink</groupId>
83+
<artifactId>flink-cdc-pipeline-connector-kafka-tis</artifactId>
84+
<version>${flink.cdc.version}</version>
85+
</dependency>
86+
8187

8288
<dependency>
8389
<groupId>com.qlangtech.tis.plugins</groupId>
@@ -93,17 +99,7 @@
9399
</dependency>
94100

95101

96-
<!-- <dependency>-->
97-
<!-- <groupId>org.apache.flink</groupId>-->
98-
<!-- <artifactId>flink-connector-kafka</artifactId>-->
99-
<!-- <version>${flink.version}</version>-->
100-
<!-- <exclusions>-->
101-
<!-- <exclusion>-->
102-
<!-- <artifactId>flink-core</artifactId>-->
103-
<!-- <groupId>org.apache.flink</groupId>-->
104-
<!-- </exclusion>-->
105-
<!-- </exclusions>-->
106-
<!-- </dependency>-->
102+
107103

108104

109105

@@ -118,5 +114,21 @@
118114

119115
</dependencies>
120116

117+
<dependencyManagement>
118+
<!-- <dependencies>-->
119+
<!-- <dependency>-->
120+
<!-- <groupId>org.apache.flink</groupId>-->
121+
<!-- <artifactId>flink-connector-kafka</artifactId>-->
122+
<!-- <version>4.0.0</version>-->
123+
<!-- <exclusions>-->
124+
<!-- <exclusion>-->
125+
<!-- <artifactId>flink-core</artifactId>-->
126+
<!-- <groupId>org.apache.flink</groupId>-->
127+
<!-- </exclusion>-->
128+
<!-- </exclusions>-->
129+
<!-- </dependency>-->
130+
<!-- </dependencies>-->
131+
</dependencyManagement>
132+
121133

122134
</project>
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,75 @@
1+
package com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.kafka.sink;
2+
3+
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
4+
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
5+
import com.qlangtech.tis.datax.IDataxProcessor;
6+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
7+
import com.qlangtech.tis.plugins.datax.kafka.writer.DataXKafkaWriter;
8+
import com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.PipelineEventSinkFunc;
9+
import com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.PipelineFlinkCDCSinkFactory;
10+
import org.apache.flink.cdc.common.factories.Factory;
11+
import org.apache.flink.cdc.common.sink.DataSink;
12+
13+
import java.time.ZoneId;
14+
import java.util.List;
15+
import java.util.Optional;
16+
17+
/**
18+
* <a href="https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/docs/connectors/pipeline-connectors/kafka/">...</a>
19+
* @author 百岁 (baisui@qlangtech.com)
20+
* @date 2026/3/13
21+
* // @see tis-flink-pipeline-paimon-plugin:com.qlangtech.tis.plugins.incr.flink.pipeline.paimon.sink.PaimonPipelineEventSinkFunc
22+
* @see DataXKafkaWriter
23+
*/
24+
public class KafkaPipelineEventSinkFunc extends PipelineEventSinkFunc<DataXKafkaWriter> {
25+
26+
/**
27+
*
28+
* @param dataxProcessor
29+
* @param pipelineSinkFactory
30+
* @param sinkDBName
31+
* @param tabs
32+
* @param sourceFlinkColCreator
33+
* @param sinkTaskParallelism
34+
*/
35+
public KafkaPipelineEventSinkFunc(IDataxProcessor dataxProcessor
36+
, PipelineFlinkCDCSinkFactory pipelineSinkFactory //
37+
, Optional<String> sinkDBName //
38+
, List<ISelectedTab> tabs //
39+
, IFlinkColCreator<FlinkCol> sourceFlinkColCreator //
40+
, int sinkTaskParallelism) {
41+
super(dataxProcessor, pipelineSinkFactory, sinkDBName, tabs, sourceFlinkColCreator, null, sinkTaskParallelism);
42+
}
43+
44+
/**
45+
* @param context
46+
* @return
47+
* @see org.apache.flink.cdc.connectors.kafka.sink.KafkaDataSink
48+
*/
49+
@Override
50+
protected DataSink createDataSink(Factory.Context context) {
51+
// DeliveryGuarantee deliveryGuarantee,
52+
// Properties kafkaProperties,
53+
// PartitionStrategy partitionStrategy,
54+
// ZoneId zoneId,
55+
// SerializationSchema<Event> keySerialization,
56+
// SerializationSchema<Event> valueSerialization,
57+
// String topic,
58+
// boolean addTableToHeaderEnabled,
59+
// String customHeaders,
60+
// String tableMapping
61+
62+
String topic = this.writer.topic;
63+
ZoneId zoneId = this.pipelineSinkFactory.getTimeZone();
64+
// return new KafkaDataSink(catalogOpts, tableOptions, commitUser, partitionMaps, serializer, zoneId,
65+
// schemaOperatorUid);
66+
return null;
67+
}
68+
69+
@Override
70+
public Factory.Context createDataSinkContext() {
71+
return null;
72+
}
73+
74+
75+
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
package com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.kafka.sink;
2+
3+
import com.google.common.collect.Sets;
4+
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
5+
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
6+
import com.qlangtech.tis.compiler.incr.ICompileAndPackage;
7+
import com.qlangtech.tis.compiler.streamcode.CompileAndPackage;
8+
import com.qlangtech.tis.datax.IDataxProcessor;
9+
import com.qlangtech.tis.extension.TISExtension;
10+
import com.qlangtech.tis.plugin.IEndTypeGetter;
11+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
12+
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
13+
import com.qlangtech.tis.plugins.datax.kafka.writer.DataXKafkaWriter;
14+
import com.qlangtech.tis.plugins.incr.flink.cdc.pipeline.PipelineFlinkCDCSinkFactory;
15+
import com.qlangtech.tis.realtime.TabSinkFunc;
16+
import org.apache.flink.cdc.common.event.Event;
17+
18+
import java.util.List;
19+
import java.util.Optional;
20+
21+
/**
22+
* <a href="https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/docs/connectors/pipeline-connectors/kafka/">...</a>
23+
*
24+
* @author 百岁 (baisui@qlangtech.com)
25+
* @date 2026/3/13
26+
*/
27+
public class KafkaPipelineSinkFactory extends PipelineFlinkCDCSinkFactory<DataXKafkaWriter> {
28+
@Override
29+
public ICompileAndPackage getCompileAndPackageManager() {
30+
return new CompileAndPackage(Sets.newHashSet(KafkaPipelineSinkFactory.class));
31+
}
32+
33+
@Override
34+
protected TabSinkFunc<?, ?, Event> createPipelineEventSinkFunc(IDataxProcessor dataxProcessor //
35+
, DataXKafkaWriter writer, List<ISelectedTab> tabs //
36+
, IFlinkColCreator<FlinkCol> sourceFlinkColCreator, IncrStreamFactory streamFactory) {
37+
return new KafkaPipelineEventSinkFunc(dataxProcessor, this //
38+
, Optional.empty(), tabs, sourceFlinkColCreator,
39+
streamFactory.getParallelism());
40+
}
41+
42+
@TISExtension
43+
public static class DftDesc extends BasicPipelineSinkDescriptor {
44+
public DftDesc() {
45+
super();
46+
}
47+
48+
@Override
49+
protected IEndTypeGetter.EndType getTargetType() {
50+
return EndType.Kafka;
51+
}
52+
}
53+
}

tis-incr/tis-flink-chunjun-kafka-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/kafka/sink/ChujunKafkaSinkFactory.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -168,7 +168,7 @@ protected Function<FieldConf, ISerializationConverter<Map<String, Object>>> getS
168168

169169
sinkFuncRef.setSinkFactory(sinkFactory);
170170
sinkFuncRef.initialize();
171-
sinkFuncRef.setSinkCols(new TableCols(kfkTable.getCols()));
171+
sinkFuncRef.setSinkCols(new TableCols<CMeta>(kfkTable.getCols()));
172172

173173
//Objects.requireNonNull(sinkFuncRef.get(), "sinkFunc can not be null");
174174
sinkFuncRef.setParallelism(this.parallelism);

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,6 @@
2323

2424
import java.io.Serializable;
2525
import java.util.Map;
26-
2726
/**
2827
* @author: 百岁(baisui@qlangtech.com)
2928
* @create: 2022-09-26 11:24
@@ -40,6 +39,7 @@ public Map<String, FlinkCol> getColMapper() {
4039
}
4140

4241
public BiFunction getSourceDTOColValProcess(String colName) {
42+
4343
FlinkCol fcol = colMapper.get(colName);
4444
if (fcol == null) {
4545
return null;
Lines changed: 81 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,81 @@
1+
package com.qlangtech.tis.plugins.incr.flink.cdc;
2+
3+
import com.qlangtech.tis.extension.Describable;
4+
import com.qlangtech.tis.extension.Descriptor;
5+
import com.qlangtech.tis.extension.util.AbstractPropAssist;
6+
import com.qlangtech.tis.extension.util.PluginExtraProps;
7+
import com.qlangtech.tis.manage.common.Option;
8+
import org.apache.commons.lang3.EnumUtils;
9+
import org.apache.flink.cdc.common.configuration.ConfigOption;
10+
import org.apache.flink.cdc.common.configuration.description.HtmlFormatter;
11+
12+
import java.lang.reflect.Method;
13+
import java.util.List;
14+
import java.util.Map;
15+
import java.util.stream.Collectors;
16+
import java.util.stream.Stream;
17+
18+
/**
19+
* @author: 百岁(baisui@qlangtech.com)
20+
* @create: 2025-06-19 06:14
21+
**/
22+
public class FlinkCDCPropAssist<T extends Describable> extends AbstractPropAssist<T,
23+
org.apache.flink.cdc.common.configuration.ConfigOption> {
24+
25+
private static final Method getClazzMethod;
26+
27+
static {
28+
try {
29+
getClazzMethod = ConfigOption.class.getDeclaredMethod("getClazz");
30+
getClazzMethod.setAccessible(true);
31+
} catch (NoSuchMethodException e) {
32+
throw new RuntimeException(e);
33+
}
34+
}
35+
36+
public static <PLUGIN extends Describable> Options<PLUGIN,
37+
org.apache.flink.cdc.common.configuration.ConfigOption> createOpts(Descriptor<PLUGIN> descriptor) {
38+
FlinkCDCPropAssist props = new FlinkCDCPropAssist(descriptor);
39+
return props.createOptions();
40+
}
41+
42+
43+
static final HtmlFormatter formatter = new HtmlFormatter();
44+
45+
public FlinkCDCPropAssist(Descriptor<T> descriptor) {
46+
super(descriptor);
47+
}
48+
49+
@Override
50+
protected MarkdownHelperContent getDescription(ConfigOption configOption) {
51+
return new MarkdownHelperContent(PluginExtraProps.AsynPropHelp.create(formatter.format(configOption.description())));
52+
}
53+
54+
@Override
55+
protected Object getDefaultValue(ConfigOption configOption) {
56+
return configOption.defaultValue();
57+
}
58+
59+
@Override
60+
protected List<Option> getOptEnums(ConfigOption configOption) {
61+
62+
try {
63+
Class clazz = (Class) getClazzMethod.invoke(configOption);
64+
if (clazz.isEnum()) {
65+
Map enumMap = EnumUtils.getEnumMap(clazz);
66+
Stream<Option> stream = enumMap.keySet().stream().map((key) -> new Option((String) key));
67+
return stream.collect(Collectors.toList());
68+
}
69+
70+
} catch (Exception e) {
71+
throw new RuntimeException(e);
72+
}
73+
74+
return null;
75+
}
76+
77+
@Override
78+
protected String getDisplayName(ConfigOption configOption) {
79+
return configOption.key();
80+
}
81+
}

0 commit comments

Comments
 (0)