Skip to content

Commit cd6289a

Browse files
committed
enable flink-cdc pipeline mode
1 parent 91c8984 commit cd6289a

10 files changed

Lines changed: 53 additions & 18 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -126,7 +126,7 @@ public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, Target
126126

127127

128128
try {
129-
SourceChannel sourceChannel = new SourceChannel(sourceFuncs);
129+
SourceChannel sourceChannel = new SourceChannel(false, sourceFuncs);
130130
sourceChannel.setFocusTabs(tabs, dataXProcessor.getTabAlias(null), (tabName) -> DTOStream.createRowData());
131131
return sourceChannel;
132132
// return (JobExecutionResult) getConsumerHandle().consume(false, name, sourceChannel, dataXProcessor);

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,7 @@ public AsyncMsg<List<ReaderSource>> start(
8484

8585
KafkaSource<DTO> source = kafkaSourceBuilder.build();
8686
try {
87-
SourceChannel sourceChannel = new SourceChannel(
87+
SourceChannel sourceChannel = new SourceChannel(flinkCDCPipelineEnable,
8888
createKafkaSource(kafkaReader.bootstrapServers, source));
8989

9090
sourceChannel.setFocusTabs(tabs, dataXProcessor.getTabAlias(null)

tis-incr/tis-flink-cdc-kingbase-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/pglike/FlinkCDCPGLikeSourceFunction.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,7 +110,7 @@ public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, Target
110110
});
111111

112112

113-
SourceChannel sourceChannel = new SourceChannel(readerSources);
113+
SourceChannel sourceChannel = new SourceChannel(flinkCDCPipelineEnable,readerSources);
114114
// for (ISelectedTab tab : tabs) {
115115
sourceChannel.setFocusTabs(tabs, dataXProcessor.getTabAlias(null), DTOStream::createDispatched);
116116
//}

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -99,7 +99,7 @@ public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, Target
9999
, new DefaultTableNameConvert()
100100
, contextParamValsGetterMapper);
101101

102-
SourceChannel sourceChannel = new SourceChannel(
102+
SourceChannel sourceChannel = new SourceChannel(flinkCDCPipelineEnable,
103103
SourceChannel.getSourceFunction(dsFactory, tabs, (dbHost, dbs, tbs, debeziumProperties) -> {
104104
List<ReaderSource> sourceFunctions = createSourceFunctions(dsFactory, tabs, deserializationSchema);
105105
return sourceFunctions;

tis-incr/tis-flink-cdc-mysql-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/cdc/mysql/FlinkCDCMysqlSourceFunction.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -214,6 +214,7 @@ public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, Target
214214

215215

216216
SourceChannel sourceChannel = new SourceChannel(
217+
flinkCDCPipelineEnable,
217218
SourceChannel.getSourceFunction(
218219
dsFactory,
219220
tabs

tis-incr/tis-flink-cdc-oracle-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/oracle/FlinkCDCOracleSourceFunction.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -108,7 +108,7 @@ public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, Target
108108
, tablesInDB.getPhysicsTabName2LogicNameConvertor()
109109
, contextParamValsGetterMapper);
110110

111-
SourceChannel sourceChannel = new SourceChannel(
111+
SourceChannel sourceChannel = new SourceChannel(flinkCDCPipelineEnable,
112112
SourceChannel.getSourceFunction(dsFactory, tabs
113113
, (dbHost, dbs, tbs, debeziumProperties) -> {
114114
return dbs.getDbStream().map((databaseName) -> {

tis-incr/tis-flink-cdc-postgresql-plugin/src/main/java/com/qlangtech/plugins/incr/flink/cdc/pglike/FlinkCDCPGLikeSourceFunction.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -112,7 +112,7 @@ public AsyncMsg<List<ReaderSource>> start(
112112
});
113113

114114

115-
SourceChannel sourceChannel = new SourceChannel(readerSources);
115+
SourceChannel sourceChannel = new SourceChannel(flinkCDCPipelineEnable, readerSources);
116116
// for (ISelectedTab tab : tabs) {
117117
sourceChannel.setFocusTabs(tabs, dataXProcessor.getTabAlias(null), DTOStream::createDispatched);
118118
//}

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

Lines changed: 20 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -18,21 +18,37 @@
1818

1919
package com.qlangtech.tis.realtime;
2020

21+
import com.qlangtech.tis.plugin.ds.ISelectedTab;
2122
import com.qlangtech.tis.realtime.transfer.DTO;
2223
import org.apache.flink.util.OutputTag;
2324

24-
import java.util.Collections;
25+
import java.util.List;
2526
import java.util.Map;
27+
import java.util.Set;
28+
import java.util.stream.Collectors;
29+
import java.util.stream.Stream;
2630

2731
/**
32+
*
2833
* @author: 百岁(baisui@qlangtech.com)
2934
* @create: 2025-01-04 09:22
3035
**/
3136
public class DTOSourceTagProcessFunction extends SourceProcessFunction<DTO> {
3237
public static final String KEY_MERGE_ALL_TABS_IN_ONE_BUS = "merge_all_tabs_in_one_bus";
3338

39+
public static Set<String> createFocusTabs(boolean flinkCDCPipelineEnable, List<ISelectedTab> tabs) {
40+
return (flinkCDCPipelineEnable
41+
? Stream.of(DTOSourceTagProcessFunction.KEY_MERGE_ALL_TABS_IN_ONE_BUS)
42+
: tabs.stream().map((t) -> t.getName())).collect(Collectors.toSet());
43+
}
44+
3445
public static DTOSourceTagProcessFunction create(boolean flinkCDCPipelineEnable, Map<String, OutputTag<DTO>> tab2OutputTag) {
35-
return flinkCDCPipelineEnable ? new MergeAllTabsInOneBusProcessFunction() : new DTOSourceTagProcessFunction(tab2OutputTag);
46+
if (flinkCDCPipelineEnable) {
47+
if (tab2OutputTag.size() != 1 || !tab2OutputTag.containsKey(KEY_MERGE_ALL_TABS_IN_ONE_BUS)) {
48+
throw new IllegalStateException("the size of tab2OutputTag must be 1,but now is:" + String.join(",", tab2OutputTag.keySet()));
49+
}
50+
}
51+
return flinkCDCPipelineEnable ? new MergeAllTabsInOneBusProcessFunction(tab2OutputTag) : new DTOSourceTagProcessFunction(tab2OutputTag);
3652
}
3753

3854
public DTOSourceTagProcessFunction(Map<String, OutputTag<DTO>> tab2OutputTag) {
@@ -47,9 +63,8 @@ protected String getTableName(DTO record) {
4763

4864
static class MergeAllTabsInOneBusProcessFunction extends DTOSourceTagProcessFunction {
4965

50-
private MergeAllTabsInOneBusProcessFunction() {
51-
super(Collections.singletonMap(KEY_MERGE_ALL_TABS_IN_ONE_BUS, new OutputTag<DTO>(KEY_MERGE_ALL_TABS_IN_ONE_BUS) {
52-
}));
66+
public MergeAllTabsInOneBusProcessFunction(Map<String, OutputTag<DTO>> tab2OutputTag) {
67+
super(tab2OutputTag);
5368
}
5469

5570
@Override

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@
2525
import java.util.Map;
2626

2727
/**
28+
* 负责给主数据流打标,提供给后续流程 {@link com.qlangtech.tis.realtime.dto.DTOStream.DispatchedDTOStream} 实现分流
29+
*
2830
* @author: 百岁(baisui@qlangtech.com)
2931
* @create: 2021-10-27 10:38
3032
**/
@@ -35,6 +37,17 @@ public SourceProcessFunction(Map<String, OutputTag<RECORD_TYPE>> tab2OutputTag)
3537
this.tab2OutputTag = tab2OutputTag;
3638
}
3739

40+
/**
41+
* 在主流中为每个表打标签
42+
*
43+
* @param in The input value.
44+
* @param ctx A {@link Context} that allows querying the timestamp of the element and getting a
45+
* {@link org.apache.flink.streaming.api.TimerService} for registering timers and querying the time. The context is only
46+
* valid during the invocation of this method, do not store it.
47+
* @param _out The collector for returning result values.
48+
* @throws Exception
49+
* @see com.qlangtech.tis.realtime.dto.DTOStream.DispatchedDTOStream#addStream collected in it
50+
*/
3851
@Override
3952
public void processElement(RECORD_TYPE in, Context ctx, Collector<RECORD_TYPE> _out) throws Exception {
4053
//side_output: https://ci.apache.org/projects/flink/flink-docs-stable/dev/stream/side_output.html

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

Lines changed: 13 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import com.qlangtech.tis.plugin.ds.DataSourceFactory;
3333
import com.qlangtech.tis.plugin.ds.ISelectedTab;
3434
import com.qlangtech.tis.plugin.ds.TableInDB;
35+
import com.qlangtech.tis.realtime.DTOSourceTagProcessFunction;
3536
import com.qlangtech.tis.realtime.ReaderSource;
3637
import com.qlangtech.tis.realtime.dto.DTOStream;
3738
import com.qlangtech.tis.sql.parser.tuple.creator.EntityName;
@@ -59,21 +60,23 @@ public class SourceChannel implements AsyncMsg<List<ReaderSource>> {
5960
private final List<ReaderSource> sourceFunction;
6061
private Set<String> focusTabs = null;// = Sets.newHashSet();
6162
private Tab2OutputTag<DTOStream> tab2OutputTag = null;
63+
private final boolean flinkCDCPipelineEnable;
6264

6365
@Override
6466
public Tab2OutputTag<DTOStream> getTab2OutputTag() {
6567
return Objects.requireNonNull(tab2OutputTag);
6668
}
6769

68-
public SourceChannel(List<ReaderSource> sourceFunction) {
70+
public SourceChannel(boolean flinkCDCPipelineEnable, List<ReaderSource> sourceFunction) {
6971
if (CollectionUtils.isEmpty(sourceFunction)) {
7072
throw new IllegalArgumentException("param sourceFunction can not be empty");
7173
}
7274
this.sourceFunction = sourceFunction;
75+
this.flinkCDCPipelineEnable = flinkCDCPipelineEnable;
7376
}
7477

75-
public SourceChannel(ReaderSource sourceFunction) {
76-
this(Collections.singletonList(sourceFunction));
78+
public SourceChannel(boolean flinkCDCPipelineEnable, ReaderSource sourceFunction) {
79+
this(flinkCDCPipelineEnable, Collections.singletonList(sourceFunction));
7780
}
7881

7982
public static List<ReaderSource> getSourceFunction(
@@ -170,12 +173,15 @@ public void setFocusTabs(List<ISelectedTab> tabs, TableAliasMapper tabAliasMappe
170173
if (tabAliasMapper.isNull()) {
171174
throw new IllegalArgumentException("param tabAliasMapper can not be null");
172175
}
173-
this.focusTabs = tabs.stream().map((t) -> t.getName()).collect(Collectors.toSet());
176+
this.focusTabs = DTOSourceTagProcessFunction.createFocusTabs(this.flinkCDCPipelineEnable, tabs);
177+
// (this.flinkCDCPipelineEnable
178+
// ? Stream.of(DTOSourceTagProcessFunction.KEY_MERGE_ALL_TABS_IN_ONE_BUS)
179+
// : tabs.stream().map((t) -> t.getName())).collect(Collectors.toSet());
174180

175-
Map<TableAlias, DTOStream> tab2StreamMapper = tabs.stream().collect(
181+
Map<TableAlias, DTOStream> tab2StreamMapper = this.focusTabs.stream().collect(
176182
Collectors.toMap(
177-
(tab) -> (tabAliasMapper.getWithCheckNotNull(tab.getName()))
178-
, (t) -> dtoStreamCreator.apply(t.getName())));
183+
(name) -> (tabAliasMapper.getWithCheckNotNull(name))
184+
, (name) -> dtoStreamCreator.apply(name)));
179185
this.tab2OutputTag
180186
= new Tab2OutputTag<>(tab2StreamMapper);
181187
}

0 commit comments

Comments
 (0)