This is an automated email from the ASF dual-hosted git repository.
MartijnVisser 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 48ddbf1f [FLINK-40362][connector/kafka] Fix idle reader hang when
no-more-splits precedes metadata update (#291)
48ddbf1f is described below
commit 48ddbf1fd5d6d3d60ed30ec4ab0780f70f9cbc51
Author: Sylwester Lachiewicz <[email protected]>
AuthorDate: Thu Sep 10 14:23:28 2026 +0200
[FLINK-40362][connector/kafka] Fix idle reader hang when no-more-splits
precedes metadata update (#291)
An idle reader (subtask with no assigned splits) that receives the
no-more-splits
signal before the reply to GetMetadataUpdateEvent never forwards it to the
sub-readers
created by that update: the re-notification in handleSourceEvents was
guarded by
!pendingSplits.isEmpty(), which is false in exactly this case. The reader
then
returns NOTHING_AVAILABLE indefinitely and a bounded job never finishes.
Re-deliver the no-more-splits signal after processing the reader's first
metadata update
(captured via !isActivelyConsumingSplits before the flag flips). Replaying
only on the
first metadata update ensures that an active reader waiting for splits on
newly added
clusters/topics does not terminate prematurely before those splits and the
subsequent
no-more-splits signal arrive. notifyNoMoreSplits() is idempotent and only
fires when
isNoMoreSplits is true.
Verified:
-
DynamicKafkaSourceReaderTest#testIdleReaderFinishesWhenNoMoreSplitsArrivesBeforeMetadata:
fails on main (NOTHING_AVAILABLE), passes with this change (END_OF_INPUT).
-
DynamicKafkaSourceReaderTest#testActiveReaderWaitsForNewSplitsAfterMetadataChange:
passes on main, fails if re-delivery is unconditional on every metadata
update,
passes with the firstMetadataUpdate guard.
- SourceTestSuiteBase#testIdleReader through
DynamicKafkaSourceITTest$IntegrationTests
hangs on unfixed main whenever a run leaves a subtask idle.
Generated-by: Claude Fable 5
---
.../source/reader/DynamicKafkaSourceReader.java | 16 ++++--
.../reader/DynamicKafkaSourceReaderTest.java | 61 ++++++++++++++++++++++
2 files changed, 74 insertions(+), 3 deletions(-)
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
index 2a6bd266..0d2c9c73 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
@@ -378,6 +378,10 @@ public class DynamicKafkaSourceReader<T> implements
SourceReader<T, DynamicKafka
clustersProperties.putAll(newClustersProperties);
}
+ // Captured before the flag flips below, so the no-more-splits replay
can tell the
+ // reader's first metadata update from a later metadata change.
+ final boolean firstMetadataUpdate = !isActivelyConsumingSplits;
+
// finally mark the reader as active, if not already and add pending
splits
if (!isActivelyConsumingSplits) {
isActivelyConsumingSplits = true;
@@ -399,9 +403,15 @@ public class DynamicKafkaSourceReader<T> implements
SourceReader<T, DynamicKafka
addSplits(validPendingSplits);
pendingSplits.clear();
- if (isNoMoreSplits) {
- notifyNoMoreSplits();
- }
+ }
+
+ // Replay only on the first metadata update. On a later metadata
change the reader must
+ // wait for the enumerator to signal again after the new assignments
(that re-signal is
+ // currently missing for a sub-enumerator recreated with partitions
already assigned; see
+ // FLINK-31006), so replaying here would finish an active reader
before the new topic's
+ // splits arrive.
+ if (isNoMoreSplits && firstMetadataUpdate) {
+ notifyNoMoreSplits();
}
// Releasing the last split output can expose the runtime output as
idle immediately. Keep
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
index ebffa3e9..1ef6f23b 100644
---
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
@@ -370,6 +370,67 @@ public class DynamicKafkaSourceReaderTest extends
SourceReaderTestBase<DynamicKa
}
}
+ @Test
+ void testIdleReaderFinishesWhenNoMoreSplitsArrivesBeforeMetadata() throws
Exception {
+ TestingReaderContext context = new TestingReaderContext();
+ try (DynamicKafkaSourceReader<Integer> reader =
createReaderWithoutStart(context)) {
+ TrackingReaderOutput<Integer> readerOutput = new
TrackingReaderOutput<>();
+ reader.start();
+
+ // The enumerator signals no more splits before the reader
receives the metadata
+ // update event, e.g. when an idle reader registers after all
bounded
+ // sub-enumerators have already finished split discovery.
+ reader.notifyNoMoreSplits();
+
+ MetadataUpdateEvent metadata =
+ DynamicKafkaSourceTestHelper.getMetadataUpdateEvent(TOPIC);
+ CompletableFuture<Void> availableBeforeMetadata =
reader.isAvailable();
+ reader.handleSourceEvents(metadata);
+
+ assertThat(availableBeforeMetadata)
+ .as("the metadata update must wake up a task parked on the
earlier future")
+ .isDone();
+ assertThat(reader.pollNext(readerOutput))
+ .as(
+ "idle reader must reach END_OF_INPUT even when
no-more-splits precedes the metadata update event")
+ .isEqualTo(InputStatus.END_OF_INPUT);
+ }
+ }
+
+ @Test
+ void testActiveReaderWaitsForNewSplitsAfterMetadataChange() throws
Exception {
+ TestingReaderContext context = new TestingReaderContext();
+ try (DynamicKafkaSourceReader<Integer> reader =
createReaderWithoutStart(context)) {
+ TrackingReaderOutput<Integer> readerOutput = new
TrackingReaderOutput<>();
+ reader.start();
+
+ // First metadata update: only cluster 0 is known, so the reader
goes active with a
+ // single sub-reader.
+ KafkaStream clusterZeroOnly =
DynamicKafkaSourceTestHelper.getKafkaStream(TOPIC);
+ clusterZeroOnly.getClusterMetadataMap().remove(kafkaClusterId1);
+ reader.handleSourceEvents(
+ new
MetadataUpdateEvent(Collections.singleton(clusterZeroOnly)));
+
+ // The enumerator finished discovery for that metadata epoch.
+ reader.notifyNoMoreSplits();
+
+ // Cluster 1 appears. The reader holds no splits, so the metadata
change recreates every
+ // sub-reader and only the reader-level flag still remembers the
earlier signal. Splits
+ // and a fresh no-more-splits signal for the new metadata arrive
after this event.
+
reader.handleSourceEvents(DynamicKafkaSourceTestHelper.getMetadataUpdateEvent(TOPIC));
+
+ assertThat(reader.pollNext(readerOutput))
+ .as(
+ "reader must not finish on the stale
no-more-splits signal while a newly added cluster still awaits its splits")
+ .isEqualTo(InputStatus.NOTHING_AVAILABLE);
+
+ reader.notifyNoMoreSplits();
+ assertThat(reader.pollNext(readerOutput))
+ .as("reader finishes once the enumerator signals again for
the new metadata")
+ .isEqualTo(InputStatus.END_OF_INPUT);
+ }
+ }
+
@Test
void testAvailabilityFutureUpdates() throws Exception {
TestingReaderContext context = new TestingReaderContext();