Skip to content

Commit db2dda2

Browse files
committed
upgrade flink-cdc version to v3.4.0,and add support to flink-cdc pipeline synchronize mode
1 parent 98e4d7c commit db2dda2

59 files changed

Lines changed: 711 additions & 604 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -228,7 +228,7 @@
228228
<testcontainers.version>1.18.3</testcontainers.version>
229229
<powerjob.version>4.3.6</powerjob.version>
230230
<mariadb-java-client.version>3.4.1</mariadb-java-client.version>
231-
<flink.cdc.version>3.1.0</flink.cdc.version>
231+
<flink.cdc.version>3.4.0</flink.cdc.version>
232232
</properties>
233233

234234

tis-asyncmsg-rocketmq-plugin/src/main/java/com/qlangtech/async/message/client/consumer/BaseConsumerListener.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -66,8 +66,8 @@ public String getTopic() {
6666
return this.topic;
6767
}
6868

69-
@Override
70-
public IConsumerHandle getConsumerHandle() {
71-
return this.consumerHandle;
72-
}
69+
// @Override
70+
// public IConsumerHandle getConsumerHandle() {
71+
// return this.consumerHandle;
72+
// }
7373
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,9 @@ public void setKey(KeyedPluginStore.Key key) {
8181

8282
public FileSystemFactory getFs() {
8383
if (fileSystem == null) {
84+
if (StringUtils.isEmpty(this.fsName)) {
85+
throw new IllegalStateException("prop field:" + this.fsName + " can not be empty");
86+
}
8487
this.fileSystem = FileSystemFactory.getFsFactory(fsName);
8588
}
8689
Objects.requireNonNull(this.fileSystem, "fileSystem has not be initialized");

tis-incr/pom.xml

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

9292
<module>tis-flink-chunjun-kingbase-plugin</module>
9393
<module>tis-flink-cdc-kingbase-plugin</module>
94-
<module>tis-flink-pipeline-paimon-plugin</module>
94+
9595
</modules>
9696

9797
<dependencyManagement>

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/source/ChunjunSourceFactory.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -50,10 +50,10 @@ public IConsumerHandle getConsumerHandle() {
5050
}
5151

5252

53-
@Override
54-
public void setConsumerHandle(IConsumerHandle consumerHandle) {
55-
this.consumerHandle = consumerHandle;
56-
}
53+
// @Override
54+
// public void setConsumerHandle(IConsumerHandle consumerHandle) {
55+
// this.consumerHandle = consumerHandle;
56+
// }
5757

5858
@Override
5959
public Descriptor<MQListenerFactory> getDescriptor() {

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/source/ChunjunSourceFunction.java

Lines changed: 9 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import com.google.common.collect.Lists;
2828
import com.google.common.collect.Maps;
2929
import com.qlangtech.plugins.incr.flink.cdc.SourceChannel;
30+
import com.qlangtech.tis.async.message.client.consumer.AsyncMsg;
3031
import com.qlangtech.tis.async.message.client.consumer.IConsumerHandle;
3132
import com.qlangtech.tis.async.message.client.consumer.IMQListener;
3233
import com.qlangtech.tis.async.message.client.consumer.MQConsumeException;
@@ -62,18 +63,18 @@
6263
* @create: 2022-08-10 11:59
6364
**/
6465
public abstract class ChunjunSourceFunction
65-
implements IMQListener<JobExecutionResult> {
66+
implements IMQListener<List<ReaderSource>> {
6667
final ChunjunSourceFactory sourceFactory;
6768

6869
public ChunjunSourceFunction(ChunjunSourceFactory sourceFactory) {
6970
this.sourceFactory = sourceFactory;
7071
}
7172

7273

73-
@Override
74-
public IConsumerHandle getConsumerHandle() {
75-
return this.sourceFactory.getConsumerHandle();
76-
}
74+
// @Override
75+
// public IConsumerHandle getConsumerHandle() {
76+
// return this.sourceFactory.getConsumerHandle();
77+
// }
7778

7879
private SourceFunction<RowData> createSourceFunction(
7980
String sourceTabName, SyncConf conf, BasicDataSourceFactory sourceFactory, BasicDataXRdbmsReader reader) {
@@ -93,7 +94,7 @@ private SourceFunction<RowData> createSourceFunction(
9394

9495

9596
@Override
96-
public JobExecutionResult start(TargetResName name, IDataxReader dataSource
97+
public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, TargetResName name, IDataxReader dataSource
9798
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
9899
Objects.requireNonNull(dataXProcessor, "dataXProcessor can not be null");
99100
BasicDataXRdbmsReader reader = (BasicDataXRdbmsReader) dataSource;
@@ -127,7 +128,8 @@ public JobExecutionResult start(TargetResName name, IDataxReader dataSource
127128
try {
128129
SourceChannel sourceChannel = new SourceChannel(sourceFuncs);
129130
sourceChannel.setFocusTabs(tabs, dataXProcessor.getTabAlias(null), (tabName) -> DTOStream.createRowData());
130-
return (JobExecutionResult) getConsumerHandle().consume(name, sourceChannel, dataXProcessor);
131+
return sourceChannel;
132+
// return (JobExecutionResult) getConsumerHandle().consume(false, name, sourceChannel, dataXProcessor);
131133
} catch (Exception e) {
132134
throw new MQConsumeException(e.getMessage(), e);
133135
}

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -134,7 +134,8 @@ public DataStreamSink<?> consumeDataStream(DataStream<Tuple2<Boolean, Row>> data
134134
return rConverter.toInternal(row.f1);
135135
}));
136136

137-
return this.rowDataSinkFunc.add2Sink(rowData);
137+
this.rowDataSinkFunc.add2Sink(rowData);
138+
return null;
138139
}
139140

140141
@Override

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

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -164,8 +164,8 @@ public static List<Option> getSupportSemantic() {
164164

165165

166166
@Override
167-
public Map<TableAlias, TabSinkFunc<RowData>> createSinkFunction(IDataxProcessor dataxProcessor, IFlinkColCreator flinkColCreator) {
168-
Map<TableAlias, TabSinkFunc<RowData>> sinkFuncs = Maps.newHashMap();
167+
public Map<TableAlias, TabSinkFunc<?, ?, RowData>> createSinkFunction(IDataxProcessor dataxProcessor, IFlinkColCreator flinkColCreator) {
168+
Map<TableAlias, TabSinkFunc<?, ?, RowData>> sinkFuncs = Maps.newHashMap();
169169

170170

171171
TableAliasMapper selectedTabs = dataxProcessor.getTabAlias(null);
@@ -237,13 +237,14 @@ public RowDataSinkFunc createRowDataSinkFunc(IDataxProcessor dataxProcessor
237237
, IPluginContext.namedContext(dataxProcessor.identityValue())
238238
, tab
239239
, sourceFlinkColCreator
240+
/** sinkColsMeta*/
240241
, AbstractRowDataMapper.getAllTabColsMeta(
241242
Objects.requireNonNull(sinkFunc.tableCols, "tabCols can not be null").getCols())
242243
, supportUpsetDML()
243244
, filterRowKinds
244245
, this.parallelism
245246
, RowDataSinkFunc.createTransformerRules(dataxProcessor.identityValue()
246-
, tabName
247+
// , tabName
247248
, tab
248249
, Objects.requireNonNull(sourceFlinkColCreator, "sourceFlinkColCreator can not be null")));
249250
}

tis-incr/tis-flink-cdc-kafka-plugin/src/main/java/com/qlangtech/tis/plugin/kafka/consumer/FlinkKafkaFunction.java

Lines changed: 6 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,8 @@
1919
package com.qlangtech.tis.plugin.kafka.consumer;
2020

2121
import com.google.common.collect.Maps;
22-
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
2322
import com.qlangtech.plugins.incr.flink.cdc.SourceChannel;
24-
import com.qlangtech.tis.async.message.client.consumer.IConsumerHandle;
25-
import com.qlangtech.tis.async.message.client.consumer.IFlinkColCreator;
23+
import com.qlangtech.tis.async.message.client.consumer.AsyncMsg;
2624
import com.qlangtech.tis.async.message.client.consumer.IMQListener;
2725
import com.qlangtech.tis.async.message.client.consumer.MQConsumeException;
2826
import com.qlangtech.tis.coredefine.module.action.TargetResName;
@@ -39,7 +37,6 @@
3937
import com.qlangtech.tis.realtime.dto.DTOStream;
4038
import com.qlangtech.tis.realtime.transfer.DTO;
4139
import com.qlangtech.tis.util.IPluginContext;
42-
import org.apache.flink.api.common.JobExecutionResult;
4340
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
4441
import org.apache.flink.api.connector.source.Source;
4542
import org.apache.flink.connector.kafka.source.KafkaSource;
@@ -60,22 +57,17 @@
6057
* @author: 百岁(baisui@qlangtech.com)
6158
* @create: 2024-12-27
6259
*/
63-
public class FlinkKafkaFunction implements IMQListener<JobExecutionResult> {
60+
public class FlinkKafkaFunction implements IMQListener<List<ReaderSource>> {
6461

6562
private final KafkaMQListenerFactory sourceFactory;
6663

67-
6864
public FlinkKafkaFunction(KafkaMQListenerFactory sourceFactory) {
6965
this.sourceFactory = sourceFactory;
7066
}
7167

7268
@Override
73-
public IConsumerHandle getConsumerHandle() {
74-
return this.sourceFactory.getConsumerHander();
75-
}
76-
77-
@Override
78-
public JobExecutionResult start(TargetResName dataxName, IDataxReader dataSource
69+
public AsyncMsg<List<ReaderSource>> start(
70+
boolean flinkCDCPipelineEnable, TargetResName dataxName, IDataxReader dataSource
7971
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
8072
DataXKafkaReader kafkaReader = (DataXKafkaReader) dataSource;
8173

@@ -98,7 +90,8 @@ public JobExecutionResult start(TargetResName dataxName, IDataxReader dataSource
9890
sourceChannel.setFocusTabs(tabs, dataXProcessor.getTabAlias(null)
9991
, (tabName) -> createDispatched(tabName, sourceFactory.independentBinLogMonitor));
10092

101-
return (JobExecutionResult) this.getConsumerHandle().consume(dataxName, sourceChannel, dataXProcessor);
93+
return sourceChannel;
94+
//return (JobExecutionResult) this.getConsumerHandle().consume(dataxName, sourceChannel, dataXProcessor);
10295
} catch (Exception e) {
10396
throw new RuntimeException(e);
10497
}

tis-incr/tis-flink-cdc-kafka-plugin/src/main/java/com/qlangtech/tis/plugin/kafka/consumer/KafkaMQListenerFactory.java

Lines changed: 14 additions & 91 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
// */
1919
package com.qlangtech.tis.plugin.kafka.consumer;
2020

21+
import com.qlangtech.plugins.incr.flink.cdc.FlinkCDCPipelineEventProcess;
2122
import com.qlangtech.plugins.incr.flink.cdc.FlinkCol;
2223
import com.qlangtech.plugins.incr.flink.cdc.RowFieldGetterFactory;
2324
import com.qlangtech.tis.annotation.Public;
@@ -54,7 +55,6 @@
5455
@Public
5556
public class KafkaMQListenerFactory extends MQListenerFactory implements KeyedPluginStore.IPluginKeyAware {
5657

57-
private transient IConsumerHandle consumerHandle;
5858

5959
/**
6060
* binlog监听在独立的slot中执行
@@ -66,51 +66,13 @@ public class KafkaMQListenerFactory extends MQListenerFactory implements KeyedPl
6666
public StartOffset startOffset;
6767
private transient String dataXName;
6868

69-
// @FormField(ordinal = 3, validate = {Validator.require})
70-
// public KafkaSubscriptionMethod subscription;
71-
72-
// @FormField(validate = {Validator.require, Validator.identity}, ordinal = 2)
73-
// public String consumeName;
74-
//
75-
// @FormField(validate = {Validator.require, Validator.identity}, ordinal = 0)
76-
// public String mqTopic;
77-
//
78-
// @FormField(validate = {Validator.require, Validator.host}, ordinal = 3)
79-
// public String namesrvAddr;
80-
//
81-
//
82-
// public String getConsumeName() {
83-
// return consumeName;
84-
// }
85-
86-
8769
@Override
8870
public void setKey(Key key) {
8971
this.dataXName = key.keyVal.getVal();
9072
}
9173

9274
@Override
9375
public IMQListener create() {
94-
// if (StringUtils.isEmpty(this.consumeName)) {
95-
// throw new IllegalStateException("prop consumeName can not be null");
96-
// }
97-
// if (StringUtils.isEmpty(this.mqTopic)) {
98-
// throw new IllegalStateException("prop mqTopic can not be null");
99-
// }
100-
// if (StringUtils.isEmpty(this.namesrvAddr)) {
101-
// throw new IllegalStateException("prop namesrvAddr can not be null");
102-
// }
103-
// if (deserialize == null) {
104-
// throw new IllegalStateException("prop deserialize can not be null");
105-
// }
106-
// ConsumerListenerForRm rmListener = new ConsumerListenerForRm();
107-
// rmListener.setConsumerGroup(this.consumeName);
108-
// rmListener.setTopic(this.mqTopic);
109-
// rmListener.setNamesrvAddr(this.namesrvAddr);
110-
// rmListener.setDeserialize(deserialize);
111-
// return rmListener;
112-
113-
11476
return new FlinkKafkaFunction(this);
11577
}
11678

@@ -150,6 +112,8 @@ public FlinkCol timestampType(DataType type) {
150112
return new FlinkCol(meta, type, new AtomicDataType(new TimestampType(nullable, 3)) //DataTypes.TIMESTAMP(3)
151113
, new KafkaTimestampValueDTOConvert(format)
152114
, new KafkaDatetimeValueDTOConvert(format)
115+
, new FlinkCDCPipelineEventProcess(
116+
org.apache.flink.cdc.common.types.DataTypes.TIMESTAMP(), new KafkaFlinkCDCPipelineTimestampValueDTOConvert(format))
153117
, new RowFieldGetterFactory.TimestampGetter(meta.getName(), colIndex));
154118

155119
}
@@ -165,6 +129,17 @@ public Object apply(Object timestamp) {
165129
}
166130
}
167131

132+
static class KafkaFlinkCDCPipelineTimestampValueDTOConvert extends KafkaDatetimeValueDTOConvert {
133+
public KafkaFlinkCDCPipelineTimestampValueDTOConvert(FormatFactory format) {
134+
super(format);
135+
}
136+
137+
@Override
138+
public Object apply(Object timestamp) {
139+
return org.apache.flink.cdc.common.data.TimestampData.fromLocalDateTime((LocalDateTime) super.apply(timestamp));
140+
}
141+
}
142+
168143
/**
169144
* @see TimestampDataConvert
170145
*/
@@ -186,27 +161,6 @@ public Object apply(Object o) {
186161
return timestamp;
187162
}
188163
}
189-
190-
// @Override
191-
// public FlinkCol varcharType(DataType type) {
192-
// FlinkCol flinkCol = super.varcharType(type);
193-
// return flinkCol.setSourceDTOColValProcess(new MySQLStringValueDTOConvert());
194-
// }
195-
//
196-
// @Override
197-
// public FlinkCol blobType(DataType type) {
198-
// FlinkCol flinkCol = super.blobType(type);
199-
// return flinkCol.setSourceDTOColValProcess(new MySQLBinaryRawValueDTOConvert());
200-
// }
201-
}
202-
203-
public IConsumerHandle getConsumerHander() {
204-
return Objects.requireNonNull(consumerHandle, "consumerHandle can not be null");
205-
}
206-
207-
@Override
208-
public void setConsumerHandle(IConsumerHandle consumerHandle) {
209-
this.consumerHandle = Objects.requireNonNull(consumerHandle, "consumerHandle can not be null");
210164
}
211165

212166
@TISExtension(ordinal = 0)
@@ -226,36 +180,5 @@ public PluginVender getVender() {
226180
public EndType getEndType() {
227181
return EndType.Kafka;
228182
}
229-
230-
231-
// private static final Pattern spec_pattern = Pattern.compile("[\\da-z_]+");
232-
233-
// private static final Pattern host_pattern = Pattern.compile("[\\da-z]{1}[\\da-z.:]+");
234-
235-
// public static final String MSG_HOST_IP_ERROR = "必须由IP、HOST及端口号组成";
236-
237-
// public static final String MSG_DIGITAL_Alpha_CHARACTER_ERROR = "必须由数字、小写字母、下划线组成";
238-
239-
// public boolean validateNamesrvAddr(IFieldErrorHandler msgHandler, Context context, String fieldName, String value) {
240-
// Matcher matcher = host_pattern.matcher(value);
241-
// if (!matcher.matches()) {
242-
// msgHandler.addFieldError(context, fieldName, MSG_HOST_IP_ERROR);
243-
// return false;
244-
// }
245-
// return true;
246-
// }
247-
248-
// public boolean validateMqTopic(IFieldErrorHandler msgHandler, Context context, String fieldName, String value) {
249-
// return validateConsumeName(msgHandler, context, fieldName, value);
250-
// }
251-
//
252-
// public boolean validateConsumeName(IFieldErrorHandler msgHandler, Context context, String fieldName, String value) {
253-
// Matcher matcher = spec_pattern.matcher(value);
254-
// if (!matcher.matches()) {
255-
// msgHandler.addFieldError(context, fieldName, MSG_DIGITAL_Alpha_CHARACTER_ERROR);
256-
// return false;
257-
// }
258-
// return true;
259-
// }
260183
}
261184
}

0 commit comments

Comments
 (0)