Skip to content

Commit 77bb05d

Browse files
committed
enable flink-cdc pipeline mode
1 parent cdea685 commit 77bb05d

2 files changed

Lines changed: 9 additions & 2 deletions

File tree

tis-incr/tis-chunjun-base-plugin/src/main/java/com/qlangtech/tis/plugins/incr/flink/chunjun/script/ChunjunSqlType.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import com.qlangtech.tis.plugin.ds.IColMetaGetter;
2828
import com.qlangtech.tis.plugins.incr.flink.connector.ChunjunSinkFactory;
2929
import com.qlangtech.tis.plugins.incr.flink.connector.streamscript.BasicFlinkStreamScriptCreator;
30+
import com.qlangtech.tis.sql.parser.tuple.creator.AdapterStreamTemplateData;
3031
import com.qlangtech.tis.sql.parser.tuple.creator.IStreamIncrGenerateStrategy;
3132

3233
import java.util.List;
@@ -68,7 +69,7 @@ public IStreamTemplateData decorateMergeData(IStreamTemplateData mergeData) {
6869
}
6970
}
7071

71-
public static class ChunjunTemplateData extends IStreamIncrGenerateStrategy.AdapterStreamTemplateData {
72+
public static class ChunjunTemplateData extends AdapterStreamTemplateData {
7273
private IStreamTableMeataCreator.ISinkStreamMetaCreator sinkStreamMetaGetter;
7374
private final IEndTypeGetter.EndType endType;
7475

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

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

1919
package com.qlangtech.tis.realtime;
2020

21+
import com.qlangtech.tis.datax.TableAlias;
2122
import com.qlangtech.tis.plugin.ds.ISelectedTab;
2223
import com.qlangtech.tis.realtime.transfer.DTO;
2324
import org.apache.flink.util.OutputTag;
@@ -29,13 +30,18 @@
2930
import java.util.stream.Stream;
3031

3132
/**
32-
*
3333
* @author: 百岁(baisui@qlangtech.com)
3434
* @create: 2025-01-04 09:22
3535
**/
3636
public class DTOSourceTagProcessFunction extends SourceProcessFunction<DTO> {
3737
public static final String KEY_MERGE_ALL_TABS_IN_ONE_BUS = "merge_all_tabs_in_one_bus";
3838

39+
public static TableAlias createAllMergeTableAlias() {
40+
return TableAlias.create(DTOSourceTagProcessFunction.KEY_MERGE_ALL_TABS_IN_ONE_BUS
41+
, DTOSourceTagProcessFunction.KEY_MERGE_ALL_TABS_IN_ONE_BUS);
42+
}
43+
44+
3945
public static Set<String> createFocusTabs(boolean flinkCDCPipelineEnable, List<ISelectedTab> tabs) {
4046
return (flinkCDCPipelineEnable
4147
? Stream.of(DTOSourceTagProcessFunction.KEY_MERGE_ALL_TABS_IN_ONE_BUS)

0 commit comments

Comments
 (0)