Skip to content

Commit 5f19953

Browse files
committed
add log
Signed-off-by: Pei Yu <125331682@qq.com>
1 parent 18b1971 commit 5f19953

1 file changed

Lines changed: 11 additions & 1 deletion

File tree

flink-cdc-flink1-compat/src/main/java/org/apache/flink/connector/base/source/reader/SingleThreadMultiplexSourceReaderBaseAdapter.java

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,19 +6,29 @@
66
import org.apache.flink.configuration.Configuration;
77
import org.apache.flink.connector.base.source.reader.fetcher.SingleThreadFetcherManager;
88

9+
import org.slf4j.Logger;
10+
import org.slf4j.LoggerFactory;
11+
912
import javax.annotation.Nullable;
1013

1114
public abstract class SingleThreadMultiplexSourceReaderBaseAdapter<
1215
E, T, SplitT extends SourceSplit, SplitStateT>
1316
extends SingleThreadMultiplexSourceReaderBase<E, T, SplitT, SplitStateT> {
1417

18+
private static final Logger LOG =
19+
LoggerFactory.getLogger(SingleThreadMultiplexSourceReaderBaseAdapter.class);
20+
1521
public SingleThreadMultiplexSourceReaderBaseAdapter(
1622
SingleThreadFetcherManager<E, SplitT> splitFetcherManager,
1723
RecordEmitter<E, T, SplitStateT> recordEmitter,
1824
@Nullable RecordEvaluator<T> eofRecordEvaluator,
1925
Configuration config,
2026
SourceReaderContext context,
21-
RateLimiterStrategy rateLimiterStrategy) {
27+
@Nullable RateLimiterStrategy rateLimiterStrategy) {
2228
super(splitFetcherManager, recordEmitter, eofRecordEvaluator, config, context);
29+
if (null != rateLimiterStrategy) {
30+
LOG.warn(
31+
"Because the runtime environment is Flink 1.x, the connector options `records.per.second` is ignored.");
32+
}
2333
}
2434
}

0 commit comments

Comments
 (0)