This is an automated email from the ASF dual-hosted git repository.
rmetzger pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
The following commit(s) were added to refs/heads/main by this push:
new 6540503d [FLINK-40367][Connectors/Kafka] Make Kafka consumer poll
timeout configurable (#292)
6540503d is described below
commit 6540503dce14a55acbf6f7a8bb29aca64efac68e
Author: Purushottam Sinha <[email protected]>
AuthorDate: Fri Aug 14 19:36:36 2026 +0530
[FLINK-40367][Connectors/Kafka] Make Kafka consumer poll timeout
configurable (#292)
KafkaPartitionSplitReader polled the consumer with a hardcoded 10 second
timeout. On idle partitions the split reader only reacts to split changes
after the ongoing poll returns, which can delay checkpoints by up to that
amount.
Introduce the poll.timeout.ms source property, defaulting to the previous
10 seconds, so the blocking duration can be tuned per job. The value is
parsed before the consumer is created and must not be negative, matching
the only bound KafkaConsumer#poll itself enforces; validating eagerly keeps
a misconfiguration from surfacing as an exception on the split fetcher
thread. A timeout of 0 is a valid non-blocking poll that returns whatever
is already buffered.
Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
.../content.zh/docs/connectors/datastream/kafka.md | 2 +
docs/content/docs/connectors/datastream/kafka.md | 3 ++
.../c0d94764-76a0-4c50-b617-70b1754c4612 | 1 +
.../connector/kafka/source/KafkaSourceOptions.java | 12 ++++++
.../source/reader/KafkaPartitionSplitReader.java | 26 +++++++++++-
.../reader/KafkaPartitionSplitReaderTest.java | 48 ++++++++++++++++++++++
6 files changed, 90 insertions(+), 2 deletions(-)
diff --git a/docs/content.zh/docs/connectors/datastream/kafka.md
b/docs/content.zh/docs/connectors/datastream/kafka.md
index a80d3371..c26d819a 100644
--- a/docs/content.zh/docs/connectors/datastream/kafka.md
+++ b/docs/content.zh/docs/connectors/datastream/kafka.md
@@ -218,6 +218,8 @@ Kafka Source 支持流式和批式两种运行模式。默认情况下,KafkaSo
请参阅下面的<a href="#dynamic-partition-discovery">动态分区检查</a>一节
- ```register.consumer.metrics``` 指定是否在 Flink 中注册 Kafka Consumer 的指标
- ```commit.offsets.on.checkpoint``` 指定是否在进行 checkpoint 时将消费位点提交至 Kafka broker
+- ```poll.timeout.ms``` 指定 Kafka Consumer 单次 poll 等待数据的最长时间(毫秒),默认为 10 秒。
+ 减小该值可以让空闲的 Source 更快地响应分片变更,但 poll 会更频繁
Kafka consumer 的配置可以参考 [Apache Kafka
文档](http://kafka.apache.org/documentation/#consumerconfigs)。
diff --git a/docs/content/docs/connectors/datastream/kafka.md
b/docs/content/docs/connectors/datastream/kafka.md
index d4b1ba22..60db801a 100644
--- a/docs/content/docs/connectors/datastream/kafka.md
+++ b/docs/content/docs/connectors/datastream/kafka.md
@@ -228,6 +228,9 @@ KafkaSource has following options for configuration:
- ```register.consumer.metrics``` specifies whether to register metrics of
KafkaConsumer in Flink
metric group
- ```commit.offsets.on.checkpoint``` specifies whether to commit consuming
offsets to Kafka brokers on checkpoint
+- ```poll.timeout.ms``` defines the maximum time in milliseconds the Kafka
consumer blocks in a
+ single poll while waiting for records, 10 seconds by default. Lowering it
makes an idle source
+ react faster to split changes, at the cost of polling more often
For configurations of KafkaConsumer, you can refer to
<a href="http://kafka.apache.org/documentation/#consumerconfigs">Apache Kafka
documentation</a>
diff --git
a/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
b/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
index afb16263..3341f09d 100644
---
a/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
+++
b/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
@@ -47,6 +47,7 @@ Method
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator
Method
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.getPendingPartitionSplitAssignment()>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaSourceEnumerator.java:0)
Method
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.getSplitOwner(org.apache.kafka.common.TopicPartition,
int)> is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaSourceEnumerator.java:0)
Method
<org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.consumer()>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaPartitionSplitReader.java:0)
+Method
<org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.getPollTimeout()>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaPartitionSplitReader.java:0)
Method
<org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.setConsumerClientRack(java.util.Properties,
java.lang.String)> is annotated with
<org.apache.flink.annotation.VisibleForTesting> in
(KafkaPartitionSplitReader.java:0)
Method
<org.apache.flink.connector.kafka.source.reader.KafkaSourceReader.getNumAliveFetchers()>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaSourceReader.java:0)
Method
<org.apache.flink.connector.kafka.source.reader.KafkaSourceReader.getOffsetsToCommit()>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaSourceReader.java:0)
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceOptions.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceOptions.java
index f96cd3ea..5928d8a2 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceOptions.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceOptions.java
@@ -58,6 +58,18 @@ public class KafkaSourceOptions {
.defaultValue(true)
.withDescription("Whether to commit consuming offset on
checkpoint.");
+ public static final ConfigOption<Long> POLL_TIMEOUT_MS =
+ ConfigOptions.key("poll.timeout.ms")
+ .longType()
+ .defaultValue(Duration.ofSeconds(10).toMillis())
+ .withDescription(
+ "The maximum time in milliseconds the Kafka
consumer blocks in a single "
+ + "poll() call while waiting for records.
Since the poll returns "
+ + "as soon as records are available, a
smaller value mainly makes "
+ + "an idle split reader react faster to
split changes, at the cost "
+ + "of polling more often. 0 polls without
blocking. Must not be "
+ + "negative.");
+
@SuppressWarnings("unchecked")
public static <T> T getOption(
Properties props, ConfigOption<?> configOption, Function<String,
T> parser) {
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReader.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReader.java
index 256a9cf1..bf8b0a97 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReader.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReader.java
@@ -64,12 +64,12 @@ import java.util.stream.Collectors;
public class KafkaPartitionSplitReader
implements SplitReader<ConsumerRecord<byte[], byte[]>,
KafkaPartitionSplit> {
private static final Logger LOG =
LoggerFactory.getLogger(KafkaPartitionSplitReader.class);
- private static final long POLL_TIMEOUT = 10000L;
private final KafkaConsumer<byte[], byte[]> consumer;
private final Map<TopicPartition, Long> stoppingOffsets;
private final String groupId;
private final int subtaskId;
+ private final Duration pollTimeout;
private final KafkaSourceReaderMetrics kafkaSourceReaderMetrics;
@@ -90,6 +90,8 @@ public class KafkaPartitionSplitReader
String rackIdSupplier) {
this.subtaskId = context.getIndexOfSubtask();
this.kafkaSourceReaderMetrics = kafkaSourceReaderMetrics;
+ // Parsed before creating the consumer so that an invalid value does
not leak a consumer.
+ this.pollTimeout = parsePollTimeout(props);
Properties consumerProps = new Properties();
consumerProps.putAll(props);
consumerProps.setProperty(ConsumerConfig.CLIENT_ID_CONFIG,
createConsumerClientId(props));
@@ -107,7 +109,7 @@ public class KafkaPartitionSplitReader
public RecordsWithSplitIds<ConsumerRecord<byte[], byte[]>> fetch() throws
IOException {
ConsumerRecords<byte[], byte[]> consumerRecords;
try {
- consumerRecords = consumer.poll(Duration.ofMillis(POLL_TIMEOUT));
+ consumerRecords = consumer.poll(pollTimeout);
} catch (WakeupException | IllegalStateException e) {
// IllegalStateException will be thrown if the consumer is not
assigned any partitions.
// This happens if all assigned partitions are invalid or empty
(starting offset >=
@@ -278,8 +280,28 @@ public class KafkaPartitionSplitReader
return consumer;
}
+ @VisibleForTesting
+ Duration getPollTimeout() {
+ return pollTimeout;
+ }
+
// --------------- private helper method ----------------------
+ private static Duration parsePollTimeout(Properties props) {
+ long pollTimeoutMs =
+ KafkaSourceOptions.getOption(
+ props, KafkaSourceOptions.POLL_TIMEOUT_MS,
Long::parseLong);
+ // KafkaConsumer#poll only rejects negative timeouts, so validate the
same bound here. Doing
+ // it eagerly turns a misconfiguration into a readable error instead
of an exception thrown
+ // from the split fetcher thread on the first poll.
+ Preconditions.checkArgument(
+ pollTimeoutMs >= 0,
+ "Property %s should not be negative, but is %s",
+ KafkaSourceOptions.POLL_TIMEOUT_MS.key(),
+ pollTimeoutMs);
+ return Duration.ofMillis(pollTimeoutMs);
+ }
+
/**
* This Method performs Null and empty Rack Id validation and sets the
rack id to the
* client.rack Consumer Config.
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReaderTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReaderTest.java
index adf240df..95502959 100644
---
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReaderTest.java
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReaderTest.java
@@ -22,6 +22,7 @@ import org.apache.flink.configuration.Configuration;
import org.apache.flink.connector.base.source.reader.RecordsWithSplitIds;
import
org.apache.flink.connector.base.source.reader.splitreader.SplitsAddition;
import org.apache.flink.connector.base.source.reader.splitreader.SplitsChange;
+import org.apache.flink.connector.kafka.source.KafkaSourceOptions;
import
org.apache.flink.connector.kafka.source.metrics.KafkaSourceReaderMetrics;
import org.apache.flink.connector.kafka.source.split.KafkaPartitionSplit;
import org.apache.flink.connector.kafka.testutils.KafkaSourceTestEnv;
@@ -53,6 +54,7 @@ import org.junit.jupiter.params.provider.NullAndEmptySource;
import org.junit.jupiter.params.provider.ValueSource;
import java.io.IOException;
+import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
@@ -377,6 +379,52 @@ public class KafkaPartitionSplitReaderTest {
Collections.singletonList(unassignedSplit));
}
+ @Test
+ public void testDefaultPollTimeout() {
+ // When the property is not set, the reader falls back to the option's
default instead of
+ // hardcoding a timeout of its own.
+ assertThat(createReader().getPollTimeout())
+
.isEqualTo(Duration.ofMillis(KafkaSourceOptions.POLL_TIMEOUT_MS.defaultValue()));
+ }
+
+ @Test
+ public void testConfiguredPollTimeout() {
+ final Properties props = new Properties();
+ props.setProperty(KafkaSourceOptions.POLL_TIMEOUT_MS.key(), "500");
+ KafkaPartitionSplitReader reader =
+ createReader(props,
UnregisteredMetricsGroup.createSourceReaderMetricGroup());
+
+ assertThat(reader.getPollTimeout()).isEqualTo(Duration.ofMillis(500));
+ }
+
+ @Test
+ public void testZeroPollTimeout() {
+ // KafkaConsumer#poll accepts a zero timeout, which returns
immediately with whatever is
+ // already buffered, so the reader must not reject it either.
+ final Properties props = new Properties();
+ props.setProperty(KafkaSourceOptions.POLL_TIMEOUT_MS.key(), "0");
+ KafkaPartitionSplitReader reader =
+ createReader(props,
UnregisteredMetricsGroup.createSourceReaderMetricGroup());
+
+ assertThat(reader.getPollTimeout()).isEqualTo(Duration.ZERO);
+ }
+
+ @Test
+ public void testNegativePollTimeoutIsRejected() {
+ final Properties props = new Properties();
+ props.setProperty(KafkaSourceOptions.POLL_TIMEOUT_MS.key(), "-1");
+ assertThatThrownBy(
+ () ->
+ createReader(
+ props,
+
UnregisteredMetricsGroup.createSourceReaderMetricGroup()))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining(
+ String.format(
+ "Property %s should not be negative, but is
-1",
+ KafkaSourceOptions.POLL_TIMEOUT_MS.key()));
+ }
+
// ------------------
private void assignSplitsAndFetchUntilFinish(KafkaPartitionSplitReader
reader, int readerId)