Skip to content

Commit 1bba664

Browse files
committed
add for paimon sink
1 parent ff5c020 commit 1bba664

33 files changed

Lines changed: 297 additions & 440 deletions

File tree

tis-incr/pom.xml

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,23 @@
9696

9797
<dependencyManagement>
9898
<dependencies>
99+
<dependency>
100+
<groupId>com.qlangtech.tis.plugins</groupId>
101+
<artifactId>tis-flink-cdc-common</artifactId>
102+
<version>${project.version}</version>
103+
<exclusions>
104+
<exclusion>
105+
<groupId>io.debezium</groupId>
106+
<artifactId>debezium-core</artifactId>
107+
</exclusion>
108+
</exclusions>
109+
</dependency>
110+
<dependency>
111+
<groupId>io.debezium</groupId>
112+
<artifactId>debezium-core</artifactId>
113+
<!--当flink-cdc 依赖的 debezium 版本变化,此版本也需要相应变化成对应的版本-->
114+
<version>${debezium-connector.version}</version>
115+
</dependency>
99116
<dependency>
100117
<groupId>com.qlangtech.tis.plugins</groupId>
101118
<artifactId>tis-flink-dependency</artifactId>
@@ -164,8 +181,6 @@
164181
</dependencies>
165182

166183

167-
168-
169184
</dependencyManagement>
170185

171186
<build>
@@ -189,7 +204,7 @@
189204
<module>tis-flink-cdc-postgresql-shade-4-debezium-connector-postgresql</module>
190205
<module>tis-flink-cdc-mysql-shade-4-debezium-connector-mysql</module>
191206
<module>tis-flink-cdc-kingbase-shade-4-debezium-connector-postgresql</module>
192-
<!-- <module>tis-flink-cdc-kingbase-shade-plugin-extends</module>-->
207+
<!-- <module>tis-flink-cdc-kingbase-shade-plugin-extends</module>-->
193208
</modules>
194209
</profile>
195210
</profiles>

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/table/ChunjunDynamicTableFactory.java

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,19 +20,16 @@
2020

2121
import com.dtstack.chunjun.connector.jdbc.table.JdbcDynamicTableFactory;
2222
import com.google.common.collect.Sets;
23-
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
24-
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
2523
import com.qlangtech.tis.datax.DataXName;
2624
import com.qlangtech.tis.datax.IDataxProcessor;
2725
import com.qlangtech.tis.datax.StoreResourceType;
2826
import com.qlangtech.tis.datax.TableAlias;
2927
import com.qlangtech.tis.datax.impl.DataxProcessor;
30-
import com.qlangtech.tis.offline.DataxUtils;
3128
import com.qlangtech.tis.plugin.IEndTypeGetter;
3229
import com.qlangtech.tis.plugin.incr.TISSinkFactory;
3330
import com.qlangtech.tis.plugins.incr.flink.chunjun.script.ChunjunSqlType;
3431
import com.qlangtech.tis.plugins.incr.flink.connector.ChunjunSinkFactory;
35-
import com.qlangtech.tis.realtime.BasicTISSinkFactory;
32+
import com.qlangtech.tis.realtime.RowDataSinkFunc;
3633
import org.apache.commons.lang3.StringUtils;
3734
import org.apache.flink.configuration.ConfigOption;
3835
import org.apache.flink.configuration.ConfigOptions;
@@ -86,7 +83,7 @@ public DynamicTableSink createDynamicTableSink(Context context) {
8683
ChunjunSinkFactory sinKFactory = (ChunjunSinkFactory) TISSinkFactory.getIncrSinKFactory(DataXName.createDataXPipeline(dataXName));
8784
IDataxProcessor dataxProcessor = DataxProcessor.load(null, dataXName);
8885

89-
BasicTISSinkFactory.RowDataSinkFunc rowDataSinkFunc = sinKFactory.createRowDataSinkFunc(dataxProcessor
86+
RowDataSinkFunc rowDataSinkFunc = sinKFactory.createRowDataSinkFunc(dataxProcessor
9087
, dataxProcessor.getTabAlias(null).getWithCheckNotNull(sourceTableName), false);
9188

9289
// 3.封装参数

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/table/ChunjunTableSinkFactory.java

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -20,19 +20,17 @@
2020

2121
import com.google.common.collect.Maps;
2222
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
23-
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
2423
import com.qlangtech.tis.datax.DataXName;
2524
import com.qlangtech.tis.datax.IDataxProcessor;
2625
import com.qlangtech.tis.datax.StoreResourceType;
2726
import com.qlangtech.tis.datax.TableAlias;
2827
import com.qlangtech.tis.datax.impl.DataxProcessor;
29-
import com.qlangtech.tis.offline.DataxUtils;
3028
import com.qlangtech.tis.plugin.IEndTypeGetter;
3129
import com.qlangtech.tis.plugin.incr.TISSinkFactory;
3230
import com.qlangtech.tis.plugins.incr.flink.cdc.impl.RowUtils;
3331
import com.qlangtech.tis.plugins.incr.flink.chunjun.script.ChunjunSqlType;
3432
import com.qlangtech.tis.plugins.incr.flink.connector.ChunjunSinkFactory;
35-
import com.qlangtech.tis.realtime.BasicTISSinkFactory;
33+
import com.qlangtech.tis.realtime.RowDataSinkFunc;
3634
import com.qlangtech.tis.realtime.dto.DTOStream;
3735
import org.apache.commons.collections.CollectionUtils;
3836
import org.apache.commons.compress.utils.Lists;
@@ -80,15 +78,15 @@ public StreamTableSink<Tuple2<Boolean, Row>> createStreamTableSink(Map<String, S
8078
ChunjunSinkFactory sinKFactory = (ChunjunSinkFactory) TISSinkFactory.getIncrSinKFactory(DataXName.createDataXPipeline(dataXName));
8179
IDataxProcessor dataxProcessor = DataxProcessor.load(null, dataXName);
8280

83-
BasicTISSinkFactory.RowDataSinkFunc rowDataSinkFunc = sinKFactory.createRowDataSinkFunc(dataxProcessor
81+
RowDataSinkFunc rowDataSinkFunc = sinKFactory.createRowDataSinkFunc(dataxProcessor
8482
, dataxProcessor.getTabAlias(null).getWithCheckNotNull(sourceTableName), false);
8583
return new ChunjunStreamTableSink(false, endType, rowDataSinkFunc);
8684
}
8785

8886

8987
public static class ChunjunStreamTableSink implements UpsertStreamTableSink<Row> {
9088

91-
private final BasicTISSinkFactory.RowDataSinkFunc rowDataSinkFunc;
89+
private final RowDataSinkFunc rowDataSinkFunc;
9290
private String[] primaryKeys;
9391
private final IEndTypeGetter.EndType endType;
9492
/**
@@ -97,7 +95,7 @@ public static class ChunjunStreamTableSink implements UpsertStreamTableSink<Row>
9795
private final boolean isAppendOnly;
9896

9997
public ChunjunStreamTableSink(boolean isAppendOnly, IEndTypeGetter.EndType endType
100-
, BasicTISSinkFactory.RowDataSinkFunc rowDataSinkFunc) {
98+
, RowDataSinkFunc rowDataSinkFunc) {
10199
this.rowDataSinkFunc = rowDataSinkFunc;
102100

103101
this.endType = endType;

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/table/TISJdbcDymaincTableSink.java

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

2121
import com.qlangtech.tis.plugin.IEndTypeGetter;
22-
import com.qlangtech.tis.realtime.BasicTISSinkFactory;
22+
import com.qlangtech.tis.realtime.RowDataSinkFunc;
2323
import org.apache.flink.table.connector.ChangelogMode;
2424
import org.apache.flink.table.connector.sink.DynamicTableSink;
2525
import org.apache.flink.table.connector.sink.SinkFunctionProvider;
@@ -31,10 +31,10 @@
3131
**/
3232
public class TISJdbcDymaincTableSink implements DynamicTableSink {
3333

34-
private final BasicTISSinkFactory.RowDataSinkFunc rowDataSinkFunc;
34+
private final RowDataSinkFunc rowDataSinkFunc;
3535
private final IEndTypeGetter.EndType endType;
3636

37-
public TISJdbcDymaincTableSink(IEndTypeGetter.EndType endType, BasicTISSinkFactory.RowDataSinkFunc rowDataSinkFunc) {
37+
public TISJdbcDymaincTableSink(IEndTypeGetter.EndType endType, RowDataSinkFunc rowDataSinkFunc) {
3838
this.rowDataSinkFunc = rowDataSinkFunc;
3939
this.endType = endType;
4040
}

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/connector/ChunjunSinkFactory.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,8 @@
7878
import com.qlangtech.tis.plugins.incr.flink.chunjun.script.ChunjunStreamScriptType;
7979
import com.qlangtech.tis.plugins.incr.flink.chunjun.sink.SinkTabPropsExtends;
8080
import com.qlangtech.tis.realtime.BasicTISSinkFactory;
81+
import com.qlangtech.tis.realtime.RowDataSinkFunc;
82+
import com.qlangtech.tis.realtime.SelectedTableTransformerRules;
8183
import com.qlangtech.tis.realtime.TabSinkFunc;
8284
import com.qlangtech.tis.realtime.transfer.DTO.EventType;
8385
import com.qlangtech.tis.runtime.module.misc.IControlMsgHandler;
@@ -243,7 +245,7 @@ public RowDataSinkFunc createRowDataSinkFunc(IDataxProcessor dataxProcessor
243245
, supportUpsetDML()
244246
, filterRowKinds
245247
, this.parallelism
246-
, RowDataSinkFunc.createTransformerRules(dataxProcessor.identityValue()
248+
, SelectedTableTransformerRules.createTransformerRules(dataxProcessor.identityValue()
247249
// , tabName
248250
, tab
249251
, Objects.requireNonNull(sourceFlinkColCreator, "sourceFlinkColCreator can not be null")));

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,12 @@
3838
<artifactId>flink-connector-debezium</artifactId>
3939
<version>${flink.cdc.version}</version>
4040
<!-- <scope>provided</scope>-->
41+
<!-- <exclusions>-->
42+
<!-- <exclusion>-->
43+
<!-- <groupId>io.debezium</groupId>-->
44+
<!-- <artifactId>debezium-core</artifactId>-->
45+
<!-- </exclusion>-->
46+
<!-- </exclusions>-->
4147
</dependency>
4248

4349
<dependency>

tis-incr/tis-flink-cdc-common/src/main/java/com/qlangtech/plugins/incr/debuzium/DebuziumPropAssist.java

Lines changed: 0 additions & 151 deletions
Original file line numberDiff line numberDiff line change
@@ -90,155 +90,4 @@ protected List<Option> getOptEnums(Field configOption) {
9090
protected String getDisplayName(Field configOption) {
9191
return configOption.displayName();
9292
}
93-
94-
// protected Options createFlinkOptions() {
95-
// return new Options(this);
96-
// }
97-
98-
// public static class Options<T extends Describable> {
99-
// private final List<FieldTriple<T>> opts = Lists.newArrayList();
100-
// private Map<String, /*** fieldname*/PropertyType> props;
101-
//
102-
// private DebuziumPropAssist<T> propsAssist;
103-
//
104-
// public Options(DebuziumPropAssist<T> propsAssist) {
105-
// this.propsAssist = propsAssist;
106-
//// this.props = getPluginFormPropertyTypes().accept(new PluginFormProperties.IVisitor() {
107-
//// @Override
108-
//// public Map<String, PropertyType> visit(RootFormProperties props) {
109-
//// return props.propertiesType;
110-
//// }
111-
//// });
112-
// }
113-
//
114-
// public void addFieldDescriptor(String fieldName, Field configOption) {
115-
// propsAssist.addFieldDescriptor(fieldName, configOption);
116-
// }
117-
//
118-
// public void addFieldDescriptor(String fieldName, Field configOption, OverwriteProps overwriteProps) {
119-
// propsAssist.addFieldDescriptor(fieldName, configOption, overwriteProps);
120-
// }
121-
//
122-
// public void add(String fieldName, TISDebuziumProp option) {
123-
// this.add(fieldName, option, null);
124-
// }
125-
//
126-
// public void add(Field option, Function<T, Object> propGetter) {
127-
// this.add(null, TISDebuziumProp.create(option), propGetter);
128-
// }
129-
//
130-
// public void add(String fieldName, TISDebuziumProp option, Function<T, Object> propGetter) {
131-
// if (StringUtils.isNotEmpty(fieldName)) {
132-
// this.addFieldDescriptor(fieldName, option.configOption, option.overwriteProp);
133-
//
134-
// }
135-
// this.opts.add(FieldTriple.of(fieldName, option.configOption, propGetter));
136-
// }
137-
//
138-
// public Map<String, PropertyType> getProps() {
139-
// if (props == null) {
140-
// this.props = propsAssist.descriptor.getPluginFormPropertyTypes().accept(new PluginFormProperties.IVisitor() {
141-
// @Override
142-
// public Map<String, PropertyType> visit(RootFormProperties props) {
143-
// return props.propertiesType;
144-
// }
145-
// });
146-
// }
147-
// return props;
148-
// }
149-
// }
150-
151-
// public static class TISDebuziumProp {
152-
// private OverwriteProps overwriteProp = new OverwriteProps();
153-
// private final Field configOption;
154-
//
155-
// boolean hasSetOverWrite = false;
156-
//
157-
// public static TISDebuziumProp create(Field configOption) {
158-
// return new TISDebuziumProp(configOption);
159-
// }
160-
//
161-
// private TISDebuziumProp(Field configOption) {
162-
// this.configOption = configOption;
163-
// }
164-
//
165-
// public TISDebuziumProp overwriteDft(Object dftVal) {
166-
// // overwriteProp.setDftVal(dftVal);
167-
// return this.setOverwriteProp(OverwriteProps.dft(dftVal));
168-
// }
169-
//
170-
// public TISDebuziumProp overwritePlaceholder(Object placeholder) {
171-
// return this.setOverwriteProp(OverwriteProps.placeholder(placeholder));
172-
// }
173-
//
174-
// public TISDebuziumProp setOverwriteProp(OverwriteProps overwriteProp) {
175-
// if (hasSetOverWrite) {
176-
// throw new IllegalStateException("overwriteProp has been setted ,can not be writen twice");
177-
// }
178-
// this.overwriteProp = overwriteProp;
179-
// this.hasSetOverWrite = true;
180-
// return this;
181-
// }
182-
// }
183-
184-
// protected void addFieldDescriptor(String fieldName, Field configOption, OverwriteProps overwriteProps) {
185-
// String desc = configOption.description();
186-
//
187-
// Object dftVal = overwriteProps.processDftVal(configOption.defaultValue());
188-
//
189-
// StringBuffer helperContent = new StringBuffer(desc);
190-
// if (overwriteProps.appendHelper.isPresent()) {
191-
// helperContent.append("\n\n").append(overwriteProps.appendHelper.get());
192-
// }
193-
//
194-
// Type targetClazz = configOption.type();
195-
// List<Option> opts = null;
196-
// switch (targetClazz) {
197-
// case LIST: {
198-
// throw new IllegalStateException("unsupported type:" + targetClazz);
199-
// }
200-
// case BOOLEAN: {
201-
// opts = Lists.newArrayList(new Option("是", true), new Option("否", false));
202-
// break;
203-
// }
204-
// case CLASS:
205-
// case PASSWORD:
206-
// case INT:
207-
// case DOUBLE:
208-
// case LONG:
209-
// case SHORT:
210-
// case STRING:
211-
// default:
212-
// // throw new IllegalStateException("unsupported type:" + targetClazz);
213-
// }
214-
//
215-
// if (configOption.recommender() instanceof EnumRecommender) {
216-
// EnumRecommender enums = (EnumRecommender) configOption.recommender();
217-
// List vals = enums.validValues(null, null);
218-
// opts = (List<Option>) vals.stream().map((e) -> new Option(String.valueOf(e))).collect(Collectors.toList());
219-
// }
220-
//
221-
//
222-
//// if (targetClazz == Duration.class) {
223-
//// if (dftVal != null) {
224-
//// dftVal = ((Duration) dftVal).getSeconds();
225-
//// }
226-
//// helperContent.append("\n\n 单位:`秒`");
227-
//// } else if (targetClazz == MemorySize.class) {
228-
//// if (dftVal != null) {
229-
//// dftVal = ((MemorySize) dftVal).getKibiBytes();
230-
//// }
231-
//// helperContent.append("\n\n 单位:`kb`");
232-
//// } else if (targetClazz.isEnum()) {
233-
//// List<Enum> enums = EnumUtils.getEnumList((Class<Enum>) targetClazz);
234-
//// opts = enums.stream().map((e) -> new Option(e.name())).collect(Collectors.toList());
235-
//// } else if (targetClazz == Boolean.class) {
236-
//// opts = Lists.newArrayList(new Option("是", true), new Option("否", false));
237-
//// }
238-
//
239-
// descriptor.addFieldDescriptor(fieldName, dftVal, configOption.displayName(), helperContent.toString()
240-
// , overwriteProps.opts.isPresent() ? overwriteProps.opts : Optional.ofNullable(opts));
241-
// }
242-
243-
24493
}

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

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,8 +50,12 @@
5050
<dependency>
5151
<groupId>com.qlangtech.tis.plugins</groupId>
5252
<artifactId>tis-flink-cdc-common</artifactId>
53-
<version>${project.version}</version>
5453
</dependency>
54+
<dependency>
55+
<groupId>io.debezium</groupId>
56+
<artifactId>debezium-core</artifactId>
57+
</dependency>
58+
5559
<!-- <dependency>-->
5660
<!-- <groupId>com.qlangtech.tis.plugins</groupId>-->
5761
<!-- <artifactId>tis-flink-cdc-postgresql-plugin</artifactId>-->

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

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,7 +110,10 @@
110110
<dependency>
111111
<groupId>com.qlangtech.tis.plugins</groupId>
112112
<artifactId>tis-flink-cdc-common</artifactId>
113-
<version>${project.version}</version>
113+
</dependency>
114+
<dependency>
115+
<groupId>io.debezium</groupId>
116+
<artifactId>debezium-core</artifactId>
114117
</dependency>
115118

116119
<!-- test dependencies on TestContainers -->

tis-incr/tis-flink-cdc-mongdb-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/mongdb/impl/MongoDBDeserializationSchema.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,6 @@
1818

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

21-
import com.alibaba.datax.common.element.Column;
2221
import com.mongodb.client.model.changestream.OperationType;
2322
import com.qlangtech.plugins.incr.flink.cdc.EventOperation;
2423
import com.qlangtech.plugins.incr.flink.cdc.ISourceValConvert;

0 commit comments

Comments
 (0)