Skip to content

Commit 11a8989

Browse files
committed
add rate limit controller setting support for TIS
1 parent 4117e8a commit 11a8989

43 files changed

Lines changed: 2100 additions & 186 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.

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

Lines changed: 0 additions & 61 deletions
This file was deleted.

tis-datax/tis-datax-local-executor-utils/src/main/java/com/qlangtech/tis/datax/DataxExecutor.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,10 +47,9 @@
4747
import com.qlangtech.tis.extension.impl.IOUtils;
4848
import com.qlangtech.tis.job.common.JobCommon;
4949
import com.qlangtech.tis.manage.common.Config;
50-
import com.qlangtech.tis.manage.common.DagTaskUtils;
50+
import com.qlangtech.tis.manage.common.TaskSoapUtils;
5151
import com.qlangtech.tis.offline.DataxUtils;
5252
import com.qlangtech.tis.order.center.IAppSourcePipelineController;
53-
import com.qlangtech.tis.datax.StoreResourceType;
5453
import com.qlangtech.tis.realtime.transfer.TableSingleDataIndexStatus;
5554
import com.qlangtech.tis.realtime.utils.NetUtils;
5655
import com.qlangtech.tis.realtime.yarn.rpc.MasterJob;
@@ -283,7 +282,7 @@ public void exec(final JarLoader uberClassLoader, DataXJobInfo jobName, IDataxPr
283282
TIS.clean(false);
284283
if (execMode == DataXJobSubmit.InstanceType.DISTRIBUTE) {
285284
try {
286-
DagTaskUtils.feedbackAsynTaskStatus(jobArgs.jobId, jobName.jobFileName, success);
285+
TaskSoapUtils.feedbackAsynTaskStatus(jobArgs.jobId, jobName.jobFileName, success);
287286
} catch (Throwable e) {
288287
logger.warn("notify exec result faild,jobId:" + jobArgs.jobId + ",jobName:" + jobName, e);
289288
}

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

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,15 @@
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.plugins.incr.flink.launch.TISFlinkCDCStreamFactory;
3031
import com.qlangtech.tis.async.message.client.consumer.AsyncMsg;
3132
import com.qlangtech.tis.async.message.client.consumer.IConsumerHandle;
3233
import com.qlangtech.tis.async.message.client.consumer.IMQListener;
3334
import com.qlangtech.tis.async.message.client.consumer.MQConsumeException;
3435
import com.qlangtech.tis.coredefine.module.action.TargetResName;
3536
import com.qlangtech.tis.datax.DataXJobInfo;
3637
import com.qlangtech.tis.datax.DataXJobSubmit;
38+
import com.qlangtech.tis.datax.DataXName;
3739
import com.qlangtech.tis.datax.IDataxProcessor;
3840
import com.qlangtech.tis.datax.IDataxReader;
3941
import com.qlangtech.tis.datax.IStreamTableMeta;
@@ -44,6 +46,7 @@
4446
import com.qlangtech.tis.plugin.ds.CMeta;
4547
import com.qlangtech.tis.plugin.ds.ISelectedTab;
4648
import com.qlangtech.tis.plugin.ds.TableInDB;
49+
import com.qlangtech.tis.plugin.incr.IncrStreamFactory;
4750
import com.qlangtech.tis.realtime.ReaderSource;
4851
import com.qlangtech.tis.realtime.dto.DTOStream;
4952
import org.apache.flink.api.common.JobExecutionResult;
@@ -65,6 +68,7 @@
6568
public abstract class ChunjunSourceFunction
6669
implements IMQListener<List<ReaderSource>> {
6770
final ChunjunSourceFactory sourceFactory;
71+
IncrStreamFactory streamFactory = new TISFlinkCDCStreamFactory();
6872

6973
public ChunjunSourceFunction(ChunjunSourceFactory sourceFactory) {
7074
this.sourceFactory = sourceFactory;
@@ -94,7 +98,7 @@ private SourceFunction<RowData> createSourceFunction(
9498

9599

96100
@Override
97-
public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, TargetResName name, IDataxReader dataSource
101+
public AsyncMsg<List<ReaderSource>> start(IncrStreamFactory streamFactory, boolean flinkCDCPipelineEnable, DataXName names, IDataxReader dataSource
98102
, List<ISelectedTab> tabs, IDataxProcessor dataXProcessor) throws MQConsumeException {
99103
Objects.requireNonNull(dataXProcessor, "dataXProcessor can not be null");
100104
BasicDataXRdbmsReader reader = (BasicDataXRdbmsReader) dataSource;
@@ -118,7 +122,7 @@ public AsyncMsg<List<ReaderSource>> start(boolean flinkCDCPipelineEnable, Target
118122
for (String physicsName : physicsNames) {
119123
SyncConf conf = createSyncConf(sourceFactory, jdbcUrl, dbName, (SelectedTab) tab, physicsName);
120124
SourceFunction<RowData> sourceFunc = createSourceFunction(tab.getName(), conf, sourceFactory, reader);
121-
sourceFuncs.add(ReaderSource.createRowDataSource(
125+
sourceFuncs.add(ReaderSource.createRowDataSource(streamFactory, names,
122126
dbHost + ":" + sourceFactory.port + "_" + dbName + "." + physicsName, tab, sourceFunc));
123127
}
124128
}
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
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.plugins.incr.flink.cdc;
20+
21+
import com.qlangtech.tis.plugin.incr.TISRateLimiter;
22+
import org.apache.flink.api.connector.source.util.ratelimit.RateLimiterStrategy;
23+
24+
/**
25+
* @author: 百岁(baisui@qlangtech.com)
26+
* @create: 2025-07-06 12:18
27+
* @see com.qlangtech.plugins.incr.flink.cdc.impl.NoRateLimiter
28+
* @see com.qlangtech.plugins.incr.flink.cdc.impl.PerSecondRateLimiter
29+
**/
30+
public abstract class BasicRateLimiter extends TISRateLimiter {
31+
32+
@Override
33+
public abstract RateLimiterStrategy getStrategy();
34+
}

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
package com.qlangtech.plugins.incr.flink.cdc;
2020

2121
import com.qlangtech.tis.realtime.transfer.DTO;
22+
import org.apache.flink.api.connector.source.util.ratelimit.RateLimiterStrategy;
2223
import org.apache.kafka.connect.data.Field;
2324

2425
import java.io.Serializable;
@@ -30,6 +31,7 @@
3031
public class DefaultSourceValConvert implements ISourceValConvert, Serializable {
3132
@Override
3233
public Object convert(DTO dto, Field field, Object val) {
34+
3335
return val;
3436
}
3537
}
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
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.plugins.incr.flink.cdc.impl;
20+
21+
import com.qlangtech.plugins.incr.flink.cdc.BasicRateLimiter;
22+
23+
import com.qlangtech.tis.extension.Descriptor;
24+
import com.qlangtech.tis.extension.TISExtension;
25+
import com.qlangtech.tis.plugin.incr.TISRateLimiter;
26+
import org.apache.flink.api.connector.source.util.ratelimit.RateLimiterStrategy;
27+
28+
/**
29+
* 不支持运行期动态设置限流值
30+
*
31+
* @author: 百岁(baisui@qlangtech.com)
32+
* @create: 2025-07-06 12:34
33+
**/
34+
public class NoRateLimiter extends BasicRateLimiter {
35+
@Override
36+
public RateLimiterStrategy getStrategy() {
37+
throw new UnsupportedOperationException();
38+
}
39+
40+
@Override
41+
public boolean supportRateLimiter() {
42+
return false;
43+
}
44+
45+
@TISExtension
46+
public static class Desc extends Descriptor<TISRateLimiter> {
47+
public Desc() {
48+
super();
49+
}
50+
51+
@Override
52+
public String getDisplayName() {
53+
return SWITCH_OFF;
54+
}
55+
}
56+
}
Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,68 @@
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.plugins.incr.flink.cdc.impl;
20+
21+
import com.alibaba.citrus.turbine.Context;
22+
import com.qlangtech.plugins.incr.flink.cdc.BasicRateLimiter;
23+
import com.qlangtech.tis.extension.Descriptor;
24+
import com.qlangtech.tis.extension.TISExtension;
25+
import com.qlangtech.tis.plugin.incr.TISRateLimiter;
26+
import com.qlangtech.tis.runtime.module.misc.IFieldErrorHandler;
27+
import org.apache.flink.api.connector.source.util.ratelimit.RateLimiterStrategy;
28+
29+
/**
30+
* @author: 百岁(baisui@qlangtech.com)
31+
* @create: 2025-07-06 12:22
32+
**/
33+
public class PerSecondRateLimiter extends BasicRateLimiter {
34+
35+
// @FormField(ordinal = 0, type = FormFieldType.INT_NUMBER, validate = {Validator.require, Validator.integer})
36+
// public Integer recordsPerSecond;
37+
38+
@Override
39+
public RateLimiterStrategy getStrategy() {
40+
return new ResettableRateLimitStrategy(Integer.MAX_VALUE);
41+
}
42+
43+
@Override
44+
public boolean supportRateLimiter() {
45+
return true;
46+
}
47+
48+
@TISExtension
49+
public static class Desc extends Descriptor<TISRateLimiter> {
50+
public Desc() {
51+
super();
52+
}
53+
54+
public boolean validateRecordsPerSecond(IFieldErrorHandler msgHandler, Context context, String fieldName, String val) {
55+
Integer limit = Integer.parseInt(val);
56+
if (limit < 100) {
57+
msgHandler.addFieldError(context, fieldName, "必须大于100");
58+
return false;
59+
}
60+
return true;
61+
}
62+
63+
@Override
64+
public String getDisplayName() {
65+
return SWITCH_ON;
66+
}
67+
}
68+
}

0 commit comments

Comments
 (0)