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();

Reply via email to