Skip to content

Commit 7cf09b2

Browse files
committed
move AbstractTabSinkFuncV1 to tis-realtime-flink
1 parent 1bba664 commit 7cf09b2

6 files changed

Lines changed: 6 additions & 29 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@
3333
<!--https://github.com/apache/flink-connector-kafka-->
3434
<!--https://nightlies.apache.org/flink/flink-docs-release-1.20/docs/connectors/datastream/kafka/-->
3535
<properties>
36-
<flink-connector-kafka.version>3.2.0-1.18</flink-connector-kafka.version>
36+
<flink-connector-kafka.version>3.4.0-1.20</flink-connector-kafka.version>
3737
</properties>
3838

3939
<dependencies>

tis-incr/tis-flink-extends/pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@
3030

3131
<artifactId>tis-flink-extends</artifactId>
3232
<properties>
33-
<flink.connector-jdbc.version>3.1.2-1.18</flink.connector-jdbc.version>
33+
<flink.connector-jdbc.version>3.3.0-1.20</flink.connector-jdbc.version>
3434
</properties>
3535

3636
<dependencies>

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

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -106,15 +106,15 @@ public JobExecutionResult consume(boolean flinkCDCPipelineEnable, TargetResName
106106
}
107107

108108
Tab2OutputTag<DTOStream> tab2OutputTag = createTab2OutputTag(asyncMsg, env, dataxName);
109-
Map<TableAlias, AbstractTabSinkFuncV1<?, ?, SINK_TRANSFER_OBJ>> sinks = createTabSinkFunc(dataXProcessor);
109+
Map<TableAlias, TabSinkFunc<?, ?, SINK_TRANSFER_OBJ>> sinks = createTabSinkFunc(dataXProcessor);
110110

111111
this.processTableStream(env, tab2OutputTag, new SinkFuncs(sinks));
112112
return executeFlinkJob(dataxName, env);
113113
}
114114

115-
protected Map<TableAlias, AbstractTabSinkFuncV1<?, ?, SINK_TRANSFER_OBJ>> createTabSinkFunc(
115+
protected Map<TableAlias, TabSinkFunc<?, ?, SINK_TRANSFER_OBJ>> createTabSinkFunc(
116116
IDataxProcessor dataXProcessor) {
117-
Map<TableAlias, AbstractTabSinkFuncV1<?, ?, SINK_TRANSFER_OBJ>> sinks
117+
Map<TableAlias, TabSinkFunc<?, ?, SINK_TRANSFER_OBJ>> sinks
118118
= this.getSinkFuncFactory().createSinkFunction(dataXProcessor, flinkColCreator);
119119
sinks.forEach((tab, func) -> {
120120
if (StringUtils.isEmpty(tab.getTo()) || StringUtils.isEmpty(tab.getFrom())) {

tis-incr/tis-flink-extends/src/main/java/com/qlangtech/tis/realtime/AbstractTabSinkFuncV1.java renamed to tis-incr/tis-realtime-flink/src/main/java/com/qlangtech/tis/realtime/AbstractTabSinkFuncV1.java

File renamed without changes.

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

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

1919
package com.qlangtech.tis.realtime;
2020

21-
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
2221
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
23-
import com.qlangtech.tis.datax.DataXName;
2422
import com.qlangtech.tis.datax.IDataxProcessor;
2523
import com.qlangtech.tis.datax.TableAlias;
26-
import com.qlangtech.tis.datax.impl.DataxProcessor;
27-
import com.qlangtech.tis.plugin.datax.transformer.RecordTransformerRules;
28-
import com.qlangtech.tis.plugin.ds.ISelectedTab;
2924
import com.qlangtech.tis.plugin.incr.TISSinkFactory;
30-
import com.qlangtech.tis.plugins.incr.flink.cdc.DTO2RowDataMapper;
31-
import com.qlangtech.tis.plugins.incr.flink.cdc.impl.RowDataTransformerMapper;
32-
import com.qlangtech.tis.realtime.dto.DTOStream;
33-
import com.qlangtech.tis.realtime.transfer.DTO;
34-
import com.qlangtech.tis.realtime.transfer.DTO.EventType;
35-
import com.qlangtech.tis.util.IPluginContext;
36-
import org.apache.commons.collections.CollectionUtils;
37-
import org.apache.flink.api.common.typeinfo.TypeInformation;
38-
import org.apache.flink.streaming.api.datastream.DataStream;
39-
import org.apache.flink.streaming.api.datastream.DataStreamSink;
40-
import org.apache.flink.streaming.api.functions.sink.SinkFunction;
41-
import org.apache.flink.table.connector.sink.SinkFunctionProvider;
42-
import org.apache.flink.table.data.RowData;
43-
import org.apache.flink.table.runtime.typeutils.InternalTypeInfo;
44-
import org.apache.flink.table.types.logical.LogicalType;
45-
import org.apache.flink.table.types.logical.RowType;
4625
import org.slf4j.Logger;
4726
import org.slf4j.LoggerFactory;
4827

49-
import java.util.List;
5028
import java.util.Map;
51-
import java.util.Optional;
5229

5330
/**
5431
* @author: 百岁(baisui@qlangtech.com)

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -106,7 +106,7 @@ protected final void processTableStream(StreamExecutionEnvironment env
106106
}
107107

108108
@Override
109-
protected Map<TableAlias, AbstractTabSinkFuncV1<?, ?, DTO>> createTabSinkFunc(IDataxProcessor dataXProcessor) {
109+
protected Map<TableAlias, TabSinkFunc<?, ?, DTO>> createTabSinkFunc(IDataxProcessor dataXProcessor) {
110110
// return super.createTabSinkFunc(dataXProcessor);
111111
return Collections.emptyMap();
112112
}

0 commit comments

Comments
 (0)