MartijnVisser commented on code in PR #309:
URL: 
https://github.com/apache/flink-connector-kafka/pull/309#discussion_r4183144959


##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java:
##########
@@ -279,6 +279,73 @@ public void 
testBoundedSourceCompletesReadersAfterActiveAndRetainedReassignment(
         }
     }
 
+    @Test
+    public void testBoundedSourceCompletesReadersAgainAfterMetadataChange() 
throws Throwable {
+        // A switchover rather than a cluster being added: the stream moves 
off cluster 0 and onto

Review Comment:
   A test that adds a topic on cluster 0 fails without the clear and passes 
with it, so a retained cluster works too. Only a cluster that gains no 
partitions can't be used here.



##########
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java:
##########
@@ -514,6 +514,14 @@ private void 
onHandleSubscribedStreamsFetch(Set<KafkaStream> fetchedKafkaStreams
         logger.info("Closing enumerators due to metadata change");
 
         closeAllEnumeratorsAndContexts();
+        // This is the point at which the previous generation of sub 
enumerators stops existing.
+        // Each one is about to be recreated and will signal no more splits 
again for its new
+        // assignments, and the reader has closed and recreated its sub 
readers, so every sub reader
+        // is back to noMoreSplitsAssignment == false. The signals recorded 
for the previous
+        // generation therefore no longer describe any live reader, and 
keeping them would make this
+        // set suppress the re-signal that a bounded job needs in order to 
finish after a metadata
+        // change.

Review Comment:
   Please cut this to one line. A cluster that gains no partitions doesn't 
signal again (FLINK-31006), and a reader only recreates its sub readers when 
the new metadata differs from what its splits cover.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to