This is an automated email from the ASF dual-hosted git repository.
baodi pushed a commit to branch branch-2.11
in repository https://gitbox.apache.org/repos/asf/pulsar.git
The following commit(s) were added to refs/heads/branch-2.11 by this push:
new 662fdaeaca0 [fix][io] Config autoCommitEnabled when it disabled
(#19499)
662fdaeaca0 is described below
commit 662fdaeaca0d1f37817047aa6f43173c8668c9ae
Author: wenbingshen <[email protected]>
AuthorDate: Tue Feb 14 21:55:18 2023 +0800
[fix][io] Config autoCommitEnabled when it disabled (#19499)
(cherry picked from commit 8b7e4ce6f1d22d20e840a8de123f810b07ca2df8)
---
.../src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java | 2 ++
1 file changed, 2 insertions(+)
diff --git
a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
index 2c4c3dd3de5..c2239748b8f 100644
---
a/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
+++
b/pulsar-io/kafka/src/main/java/org/apache/pulsar/io/kafka/KafkaAbstractSource.java
@@ -110,6 +110,8 @@ public abstract class KafkaAbstractSource<V> extends
KafkaPushSource<V> {
}
props.put(ConsumerConfig.GROUP_ID_CONFIG,
kafkaSourceConfig.getGroupId());
props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG,
String.valueOf(kafkaSourceConfig.getFetchMinBytes()));
+ props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,
+ String.valueOf(kafkaSourceConfig.isAutoCommitEnabled()));
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG,
String.valueOf(kafkaSourceConfig.getAutoCommitIntervalMs()));
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG,
String.valueOf(kafkaSourceConfig.getSessionTimeoutMs()));