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]