Savonitar commented on code in PR #305:
URL: 
https://github.com/apache/flink-connector-kafka/pull/305#discussion_r4099575585


##########
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumerator.java:
##########
@@ -368,9 +374,19 @@ private void checkPartitionChanges(Set<TopicPartition> 
fetchedPartitions, Throwa
         if (partitionChange.isEmpty()) {
             return;
         }
+        // Track in-flight partitions to avoid duplicate discovery.
+        final Set<TopicPartition> partitionsInThisBatch = new HashSet<>();
+        partitionsInThisBatch.addAll(partitionChange.getInitialPartitions());
+        partitionsInThisBatch.addAll(partitionChange.getNewPartitions());
+        partitionsBeingInitialized.addAll(partitionsInThisBatch);
+
         context.callAsync(
                 () -> initializePartitionSplits(partitionChange),
-                this::handlePartitionSplitChanges);
+                (result, error) -> {
+                    // Clear in-flight state on the success and failure paths.
+                    
partitionsBeingInitialized.removeAll(partitionsInThisBatch);

Review Comment:
   Should we also clear the in-flight state when 
`StoppableKafkaEnumContextProxy` suppresses an initialization failure?
   Wouldnt this leave the partition unassigned until the enumerator is 
recreated? 
   Could we cover this recovery sequence with a regression test and ensure the 
suppressed-failure path releases the partition?



-- 
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