Skip to content

Commit 0c8e943

Browse files
committed
extract interface IConsumerRateLimiter for incr consume rate control
1 parent f42939f commit 0c8e943

11 files changed

Lines changed: 89 additions & 22 deletions

File tree

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@
4646
import com.qlangtech.tis.plugin.ds.CMeta;
4747
import com.qlangtech.tis.plugin.ds.ISelectedTab;
4848
import com.qlangtech.tis.plugin.ds.TableInDB;
49+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
4950
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
5051
import com.qlangtech.tis.realtime.ReaderSource;
5152
import com.qlangtech.tis.realtime.dto.DTOStream;
@@ -98,7 +99,7 @@ private SourceFunction<RowData> createSourceFunction(
9899

99100

100101
@Override
101-
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, DataXName names, IDataxReader dataSource
102+
public AsyncMsg<List<ReaderSource>> start(IConsumerRateLimiter streamFactory, boolean flinkCDCPipelineEnable, DataXName names, IDataxReader dataSource
102103
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
103104
Objects.requireNonNull(dataXProcessor, "dataXProcessor can not be null");
104105
BasicDataXRdbmsReader reader = (BasicDataXRdbmsReader) dataSource;

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@
3030
import com.qlangtech.tis.plugin.datax.transformer.RecordTransformerRules;
3131
import com.qlangtech.tis.plugin.ds.ISelectedTab;
3232
import com.qlangtech.tis.plugin.ds.RunningContext;
33+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
3334
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
3435
import com.qlangtech.tis.realtime.DTOSourceTagProcessFunction;
3536
import com.qlangtech.tis.realtime.ReaderSource;
@@ -67,7 +68,7 @@ public FlinkKafkaFunction(KafkaMQListenerFactory sourceFactory) {
6768
}
6869

6970
@Override
70-
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory,
71+
public AsyncMsg<List<ReaderSource>> start(IConsumerRateLimiter streamFactory,
7172
boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
7273
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
7374
DataXKafkaReader kafkaReader = (DataXKafkaReader) dataSource;
@@ -102,7 +103,7 @@ private DTOStream<DTO> createDispatched(String table, boolean startNewChain) {
102103
return new KafkaDispatchedDTOStream(table, startNewChain);
103104
}
104105

105-
public static ReaderSource<DTO> createKafkaSource(IncrStreamFactory streamFactory, DataXName dataXName, String tokenName, Source<DTO, ?, ?> sourceFunc) {
106+
public static ReaderSource<DTO> createKafkaSource(IConsumerRateLimiter streamFactory, DataXName dataXName, String tokenName, Source<DTO, ?, ?> sourceFunc) {
106107
return new SideOutputReaderSource<DTO>(streamFactory, dataXName, tokenName) {
107108
@Override
108109
protected DataStreamSource<DTO> addAsSource(StreamExecutionEnvironment env) {

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
import com.qlangtech.tis.plugin.ds.DataSourceFactory.ISchemaSupported;
3939
import com.qlangtech.tis.plugin.ds.ISelectedTab;
4040
import com.qlangtech.tis.plugin.ds.RunningContext;
41+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
4142
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
4243
import com.qlangtech.tis.realtime.ReaderSource;
4344
import com.qlangtech.tis.realtime.dto.DTOStream;
@@ -75,7 +76,7 @@ public FlinkCDCPGLikeSourceFunction(FlinkCDCPGLikeSourceFactory sourceFactory) {
7576
// }
7677

7778
@Override
78-
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
79+
public AsyncMsg<List<ReaderSource>> start(IConsumerRateLimiter streamFactory, boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
7980
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
8081
try {
8182
BasicDataXRdbmsReader rdbmsReader = (BasicDataXRdbmsReader) dataSource;

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838
import com.qlangtech.tis.plugin.ds.ISelectedTab;
3939
import com.qlangtech.tis.plugin.ds.RunningContext;
4040
import com.qlangtech.tis.plugin.ds.mangodb.MangoDBDataSourceFactory;
41+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
4142
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
4243
import com.qlangtech.tis.plugins.incr.flink.FlinkColMapper;
4344
import com.qlangtech.tis.plugins.incr.flink.cdc.AbstractRowDataMapper;
@@ -70,7 +71,7 @@ public FlinkCDCMongoDBSourceFunction(FlinkCDCMongoDBSourceFactory sourceFactory)
7071
}
7172

7273
@Override
73-
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
74+
public AsyncMsg<List<ReaderSource>> start(IConsumerRateLimiter streamFactory, boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
7475
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
7576
try {
7677
DataXMongodbReader mongoReader = (DataXMongodbReader) dataSource;
@@ -113,7 +114,7 @@ public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory, boole
113114
}
114115
}
115116

116-
private List<ReaderSource> createSourceFunctions(IncrStreamFactory streamFactory, DataXName dataXName,
117+
private List<ReaderSource> createSourceFunctions(IConsumerRateLimiter streamFactory, DataXName dataXName,
117118
MangoDBDataSourceFactory dsFactory, List<ISelectedTab> tabs, TISDeserializationSchema deserializationSchema) {
118119
List<ReaderSource> sourceFuncs = Lists.newArrayList();
119120

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

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@
4545
import com.qlangtech.tis.plugin.ds.ISelectedTab;
4646
import com.qlangtech.tis.plugin.ds.RunningContext;
4747
import com.qlangtech.tis.plugin.ds.TableInDB;
48+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
4849
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
4950
import com.qlangtech.tis.plugins.incr.flink.FlinkColMapper;
5051
import com.qlangtech.tis.plugins.incr.flink.cdc.AbstractRowDataMapper;
@@ -189,7 +190,7 @@ public Object apply(Object o) {
189190
* @see JobExecutionResult
190191
*/
191192
@Override
192-
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
193+
public AsyncMsg<List<ReaderSource>> start(IConsumerRateLimiter streamFactory, boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
193194
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
194195
try {
195196
Objects.requireNonNull(dataXProcessor, "param dataXProcessor can not be null");
@@ -239,13 +240,13 @@ public static class MySQLReaderSourceCreator implements ReaderSourceCreator {
239240
private static final Logger logger = LoggerFactory.getLogger(MySQLReaderSourceCreator.class);
240241
private final boolean flinkCDCPipelineEnable;
241242
private final DataXName dataXName;
242-
private final IncrStreamFactory streamFactory;
243+
private final IConsumerRateLimiter streamFactory;
243244

244-
public MySQLReaderSourceCreator(DataXName dataXName, IncrStreamFactory streamFactory, BasicDataSourceFactory dsFactory, FlinkCDCMySQLSourceFactory sourceFactory) {
245+
public MySQLReaderSourceCreator(DataXName dataXName, IConsumerRateLimiter streamFactory, BasicDataSourceFactory dsFactory, FlinkCDCMySQLSourceFactory sourceFactory) {
245246
this(dataXName, streamFactory, false, dsFactory, sourceFactory, new TISDeserializationSchema());
246247
}
247248

248-
public MySQLReaderSourceCreator(DataXName dataXName, IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, BasicDataSourceFactory dsFactory
249+
public MySQLReaderSourceCreator(DataXName dataXName, IConsumerRateLimiter streamFactory, boolean flinkCDCPipelineEnable, BasicDataSourceFactory dsFactory
249250
, FlinkCDCMySQLSourceFactory sourceFactory, TISDeserializationSchema deserializationSchema) {
250251
this.dsFactory = dsFactory;
251252
this.dataXName = Objects.requireNonNull(dataXName, "dataXName can not be null");
Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
/**
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
* <p>
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
* <p>
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
19+
package com.qlangtech.tis.plugins.incr.flink.cdc.mysql;
20+
21+
import com.qlangtech.tis.TIS;
22+
import com.qlangtech.tis.plugin.ds.DataSourceFactory;
23+
import com.qlangtech.tis.plugin.ds.PostedDSProp;
24+
import org.junit.Test;
25+
26+
/**
27+
* @author: 百岁(baisui@qlangtech.com)
28+
* @create: 2025-08-21 12:12
29+
**/
30+
public class BatchUpdate {
31+
32+
@Test()
33+
public void testBatchUpdate() {
34+
DataSourceFactory dataSource = TIS.getDataBasePlugin(PostedDSProp.parse("order"));
35+
36+
dataSource.visitFirstConnection((conn) -> {
37+
int count = 0;
38+
while (true) {
39+
try {
40+
conn.execute("update orderdetail_02 set op_time=op_time+1 where last_ver = " + ((count++) % 13));
41+
Thread.sleep(1000);
42+
} catch (Exception e) {
43+
throw new RuntimeException(e);
44+
}
45+
46+
}
47+
48+
});
49+
}
50+
}

tis-incr/tis-flink-cdc-mysql-plugin/src/test/java/com/qlangtech/tis/plugins/incr/flink/cdc/mysql/TestFlinkCDCMySQLSourceFactory.java

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@
5454
import com.qlangtech.tis.plugin.ds.ISelectedTab;
5555
import com.qlangtech.tis.plugin.ds.JDBCTypes;
5656
import com.qlangtech.tis.plugin.ds.RdbmsRunningContext;
57+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
5758
import com.qlangtech.tis.plugin.incr.TISSinkFactory;
5859
import com.qlangtech.tis.plugins.incr.flink.cdc.mysql.startup.LatestStartupOptions;
5960
import com.qlangtech.tis.realtime.ReaderSource;
@@ -500,7 +501,7 @@ protected IResultRows createConsumerHandle(BasicDataXRdbmsReader dataxReader, St
500501
*/
501502
@Test
502503
public void testBinlogConsumeWithDataStreamRegisterInstaneDetailTable() throws Exception {
503-
FlinkCDCMySQLSourceFactory mysqlCDCFactory = createCDCFactory();
504+
final FlinkCDCMySQLSourceFactory mysqlCDCFactory = createCDCFactory();
504505
mysqlCDCFactory.startupOptions = new LatestStartupOptions();
505506
// final String tabName = "instancedetail";
506507

@@ -548,7 +549,7 @@ protected void manipulateAndVerfiyTableCrudProcess(String tabName, BasicDataXRdb
548549
exampleRows.add(this.parseTestRow(RowKind.INSERT, TestFlinkCDCMySQLSourceFactory.class, tabName + "/insert1.txt"));
549550

550551
Assert.assertEquals(1, exampleRows.size());
551-
imqListener.start(false, DataXName.createDataXPipeline(dataxName.getName()), dataxReader, tabs, createProcess());
552+
imqListener.start(IConsumerRateLimiter.unsuppoted(), false, DataXName.createDataXPipeline(dataxName.getName()), dataxReader, tabs, createProcess());
552553

553554
Thread.sleep(1000);
554555
CloseableIterator<Row> snapshot = consumerHandle.getRowSnapshot(tabName);
@@ -761,6 +762,11 @@ public <T extends ISelectedTab> List<T> getSelectedTabs() {
761762
throw new UnsupportedOperationException();
762763
}
763764

765+
@Override
766+
public <T extends ISelectedTab> List<T> getUnfilledSelectedTabs() {
767+
return List.of();
768+
}
769+
764770
@Override
765771
public IGroupChildTaskIterator getSubTasks(Predicate<ISelectedTab> filter) {
766772
throw new UnsupportedOperationException();

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@
3939
import com.qlangtech.tis.plugin.ds.ISelectedTab;
4040
import com.qlangtech.tis.plugin.ds.RunningContext;
4141
import com.qlangtech.tis.plugin.ds.TableInDB;
42+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
4243
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
4344
import com.qlangtech.tis.plugins.incr.flink.FlinkColMapper;
4445
import com.qlangtech.tis.plugins.incr.flink.cdc.AbstractRowDataMapper;
@@ -83,7 +84,7 @@ public FlinkCDCOracleSourceFunction(FlinkCDCOracleSourceFactory sourceFactory) {
8384

8485
@Override
8586
public AsyncMsg<List<ReaderSource>> start(
86-
IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, DataXName channalName, IDataxReader dataSource
87+
IConsumerRateLimiter streamFactory, boolean flinkCDCPipelineEnable, DataXName channalName, IDataxReader dataSource
8788
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
8889
try {
8990

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@
3737
import com.qlangtech.tis.plugin.ds.DataSourceFactory.ISchemaSupported;
3838
import com.qlangtech.tis.plugin.ds.ISelectedTab;
3939
import com.qlangtech.tis.plugin.ds.RunningContext;
40+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
4041
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
4142
import com.qlangtech.tis.realtime.ReaderSource;
4243
import com.qlangtech.tis.realtime.dto.DTOStream;
@@ -74,7 +75,7 @@ public FlinkCDCPGLikeSourceFunction(FlinkCDCPGLikeSourceFactory sourceFactory) {
7475
// }
7576

7677
@Override
77-
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory,
78+
public AsyncMsg<List<ReaderSource>> start(IConsumerRateLimiter streamFactory,
7879
boolean flinkCDCPipelineEnable, DataXName dataxName, IDataxReader dataSource
7980
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
8081
try {

tis-incr/tis-flink-chunjun-oracle-plugin/src/test/java/com/qlangtech/plugins/incr/flink/chunjun/oracle/source/TestChunjunOracleSourceFactory.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,10 +24,12 @@
2424
import com.qlangtech.plugins.incr.flink.junit.TISApplySkipFlinkClassloaderFactoryCreation;
2525
import com.qlangtech.tis.async.message.client.consumer.IMQListener;
2626
import com.qlangtech.tis.coredefine.module.action.TargetResName;
27+
import com.qlangtech.tis.datax.DataXName;
2728
import com.qlangtech.tis.plugin.datax.common.BasicDataXRdbmsReader;
2829
import com.qlangtech.tis.plugin.ds.BasicDataSourceFactory;
2930
import com.qlangtech.tis.plugin.ds.ISelectedTab;
3031
import com.qlangtech.tis.plugin.ds.oracle.OracleDSFactoryContainer;
32+
import com.qlangtech.tis.plugin.incr.IConsumerRateLimiter;
3133
import com.qlangtech.tis.plugins.incr.flink.chunjun.offset.ScanAll;
3234
import com.qlangtech.tis.realtime.ReaderSource;
3335
import org.apache.flink.api.common.JobExecutionResult;
@@ -118,7 +120,8 @@ protected List<TestRow> createExampleTestRows() throws Exception {
118120
protected void manipulateAndVerfiyTableCrudProcess(String tabName, BasicDataXRdbmsReader dataxReader
119121
, ISelectedTab tab, IResultRows consumerHandle, IMQListener<List<ReaderSource>> imqListener) throws Exception {
120122
// super.verfiyTableCrudProcess(tabName, dataxReader, tab, consumerHandle, imqListener);
121-
imqListener.start(false, dataxName, dataxReader, Collections.singletonList(tab), createProcess());
123+
imqListener.start(IConsumerRateLimiter.unsuppoted(),false
124+
, DataXName.createDataXPipeline(dataxName.getName()) , dataxReader, Collections.singletonList(tab), createProcess());
122125
CloseableIterator<Row> snapshot = consumerHandle.getRowSnapshot(tabName);
123126
waitForSnapshotStarted(snapshot);
124127
while (snapshot.hasNext()) {

0 commit comments

Comments
 (0)