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


##########
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java:
##########
@@ -394,10 +396,19 @@ private void handleNoMoreSplits() {
             }
 
             if (firstDiscoveryComplete && 
allEnumeratorsHaveSignalledNoMoreSplits) {
-                logger.info(
-                        "Signal no more splits to all readers: {}",
-                        enumContext.registeredReaders().keySet());
-                
enumContext.registeredReaders().keySet().forEach(enumContext::signalNoMoreSplits);
+                // a reader that has already been signalled may have finished, 
and
+                // NoMoreSplitsEvent is not loss tolerant: re-sending it to a 
FINISHED task fails
+                // the job. Signal each registered reader once, until it 
registers again.
+                Set<Integer> readersToSignal =
+                        new 
HashSet<>(enumContext.registeredReaders().keySet());
+                readersToSignal.removeAll(readersSignalledNoMoreSplits);
+                if (readersToSignal.isEmpty()) {

Review Comment:
   Nit: `if (readersSignalledNoMoreSplits.add(readerId))` per reader does the 
same without the temporary set and the early `return`.



##########
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java:
##########
@@ -767,6 +778,9 @@ private void addSplitsBackToClusterEnumerators(
     @Override
     public void addReader(int subtaskId) {
         logger.debug("Adding reader {}", subtaskId);
+        // this reader is (re-)registering, so a no more splits signal sent to 
a previous attempt
+        // does not count: it must be signalled again or a restarted reader 
would never finish.
+        readersSignalledNoMoreSplits.remove(subtaskId);

Review Comment:
   Resetting here is right, but not for the reason given in the description and 
the commit message. `SourceCoordinator.subtaskReset` calls `addSplitsBack` for 
every reset subtask, with an empty list for a reader that owned nothing; on 
`main` that empty call re-signals everyone (checked with a test, reader 0 ends 
at 4). Please reword both to what holds: the new attempt's registration is the 
event that makes the reader addressable again.



##########
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java:
##########
@@ -394,10 +396,19 @@ private void handleNoMoreSplits() {
             }
 
             if (firstDiscoveryComplete && 
allEnumeratorsHaveSignalledNoMoreSplits) {
-                logger.info(
-                        "Signal no more splits to all readers: {}",
-                        enumContext.registeredReaders().keySet());
-                
enumContext.registeredReaders().keySet().forEach(enumContext::signalNoMoreSplits);
+                // a reader that has already been signalled may have finished, 
and
+                // NoMoreSplitsEvent is not loss tolerant: re-sending it to a 
FINISHED task fails
+                // the job. Signal each registered reader once, until it 
registers again.

Review Comment:
   Nit: "until it registers again" is one of two resets. A metadata change 
recreates all sub-enumerators in `onHandleSubscribedStreamsFetch`, and their 
new assignments need a fresh signal; today that path never signals anyway 
(FLINK-31006), so clearing the set there changes nothing observable, but it 
keeps this set from being the thing that blocks the re-signal once 31006 is 
fixed. Fine as a follow-up if you prefer.



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java:
##########
@@ -2455,6 +2455,93 @@ private ClusterMetadata copyClusterMetadataWithOverrides(
                 baseClusterMetadata.getStoppingOffsetsInitializer());
     }
 
+    @Test
+    public void testNoMoreSplitsIsSignalledOncePerReaderRegistration() throws 
Throwable {
+        try (CountingSignalNoMoreSplitsContext context =
+                        new CountingSignalNoMoreSplitsContext(NUM_SUBTASKS);
+                DynamicKafkaSourceEnumerator enumerator = 
createBoundedEnumerator(context)) {
+            enumerator.start();
+            runAllOneTimeCallables(context);
+
+            for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
+                mockRegisterReaderAndSendReaderStartupEvent(context, 
enumerator, reader);
+                runAllOneTimeCallables(context);
+            }
+
+            for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
+                assertThat(context.getSignalCount(reader))
+                        .as("reader %s should be signalled no more splits 
once", reader)
+                        .isEqualTo(1);
+            }
+
+            // one reader fails and returns its splits. That drives 
handleNoMoreSplits() again, but
+            // the other readers may already have finished on the first 
signal, and
+            // NoMoreSplitsEvent is not loss tolerant: re-sending it to a 
FINISHED task fails the
+            // job with "An OperatorEvent from an OperatorCoordinator to a 
task was lost".
+            int failingReader = NUM_SUBTASKS - 1;
+            List<DynamicKafkaSourceSplit> splitsOfFailingReader = new 
ArrayList<>();
+            for (SplitsAssignment<DynamicKafkaSourceSplit> assignment :
+                    context.getSplitsAssignmentSequence()) {
+                List<DynamicKafkaSourceSplit> splits = 
assignment.assignment().get(failingReader);
+                if (splits != null) {
+                    splitsOfFailingReader.addAll(splits);
+                }
+            }
+            assertThat(splitsOfFailingReader).as("precondition: reader had 
splits").isNotEmpty();
+
+            enumerator.addSplitsBack(splitsOfFailingReader, failingReader);
+
+            for (int reader = 0; reader < failingReader; reader++) {
+                assertThat(context.getSignalCount(reader))
+                        .as("reader %s must not be signalled again after it 
was already told", reader)
+                        .isEqualTo(1);
+            }
+        }
+    }
+
+    private DynamicKafkaSourceEnumerator createBoundedEnumerator(

Review Comment:
   Nit: this copies `createEnumerator` with two arguments changed; a 
`Boundedness` parameter on the existing helper is enough.



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java:
##########
@@ -2455,6 +2455,93 @@ private ClusterMetadata copyClusterMetadataWithOverrides(
                 baseClusterMetadata.getStoppingOffsetsInitializer());
     }
 
+    @Test
+    public void testNoMoreSplitsIsSignalledOncePerReaderRegistration() throws 
Throwable {
+        try (CountingSignalNoMoreSplitsContext context =
+                        new CountingSignalNoMoreSplitsContext(NUM_SUBTASKS);
+                DynamicKafkaSourceEnumerator enumerator = 
createBoundedEnumerator(context)) {
+            enumerator.start();
+            runAllOneTimeCallables(context);
+
+            for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
+                mockRegisterReaderAndSendReaderStartupEvent(context, 
enumerator, reader);
+                runAllOneTimeCallables(context);
+            }
+
+            for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
+                assertThat(context.getSignalCount(reader))
+                        .as("reader %s should be signalled no more splits 
once", reader)
+                        .isEqualTo(1);
+            }
+
+            // one reader fails and returns its splits. That drives 
handleNoMoreSplits() again, but
+            // the other readers may already have finished on the first 
signal, and
+            // NoMoreSplitsEvent is not loss tolerant: re-sending it to a 
FINISHED task fails the
+            // job with "An OperatorEvent from an OperatorCoordinator to a 
task was lost".
+            int failingReader = NUM_SUBTASKS - 1;
+            List<DynamicKafkaSourceSplit> splitsOfFailingReader = new 
ArrayList<>();
+            for (SplitsAssignment<DynamicKafkaSourceSplit> assignment :
+                    context.getSplitsAssignmentSequence()) {
+                List<DynamicKafkaSourceSplit> splits = 
assignment.assignment().get(failingReader);
+                if (splits != null) {
+                    splitsOfFailingReader.addAll(splits);
+                }
+            }
+            assertThat(splitsOfFailingReader).as("precondition: reader had 
splits").isNotEmpty();
+
+            enumerator.addSplitsBack(splitsOfFailingReader, failingReader);

Review Comment:
   The second half of the contract, that a re-registered reader is signalled 
again, is not asserted, and the reset reader is still registered here, which is 
not what the coordinator produces (`executionAttemptFailed` unregisters before 
`subtaskReset`). Please unregister it first, then re-register after 
`addSplitsBack` and assert its count goes to 2 while the others stay at 1. 
Passes on this branch, fails on `main`.
   
   ```java
               context.unregisterReader(failingReader);
               enumerator.addSplitsBack(splitsOfFailingReader, failingReader);
               for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
                   assertThat(context.getSignalCount(reader)).isEqualTo(1);
               }
   
               mockRegisterReaderAndSendReaderStartupEvent(context, enumerator, 
failingReader);
               runAllOneTimeCallables(context);
               assertThat(context.getSignalCount(failingReader)).isEqualTo(2);
               for (int reader = 0; reader < failingReader; reader++) {
                   assertThat(context.getSignalCount(reader)).isEqualTo(1);
               }
   ```



##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java:
##########
@@ -2455,6 +2455,93 @@ private ClusterMetadata copyClusterMetadataWithOverrides(
                 baseClusterMetadata.getStoppingOffsetsInitializer());
     }
 
+    @Test
+    public void testNoMoreSplitsIsSignalledOncePerReaderRegistration() throws 
Throwable {
+        try (CountingSignalNoMoreSplitsContext context =
+                        new CountingSignalNoMoreSplitsContext(NUM_SUBTASKS);
+                DynamicKafkaSourceEnumerator enumerator = 
createBoundedEnumerator(context)) {
+            enumerator.start();
+            runAllOneTimeCallables(context);
+
+            for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
+                mockRegisterReaderAndSendReaderStartupEvent(context, 
enumerator, reader);
+                runAllOneTimeCallables(context);
+            }
+
+            for (int reader = 0; reader < NUM_SUBTASKS; reader++) {
+                assertThat(context.getSignalCount(reader))
+                        .as("reader %s should be signalled no more splits 
once", reader)
+                        .isEqualTo(1);
+            }
+
+            // one reader fails and returns its splits. That drives 
handleNoMoreSplits() again, but
+            // the other readers may already have finished on the first 
signal, and
+            // NoMoreSplitsEvent is not loss tolerant: re-sending it to a 
FINISHED task fails the
+            // job with "An OperatorEvent from an OperatorCoordinator to a 
task was lost".
+            int failingReader = NUM_SUBTASKS - 1;
+            List<DynamicKafkaSourceSplit> splitsOfFailingReader = new 
ArrayList<>();
+            for (SplitsAssignment<DynamicKafkaSourceSplit> assignment :
+                    context.getSplitsAssignmentSequence()) {
+                List<DynamicKafkaSourceSplit> splits = 
assignment.assignment().get(failingReader);
+                if (splits != null) {
+                    splitsOfFailingReader.addAll(splits);
+                }
+            }
+            assertThat(splitsOfFailingReader).as("precondition: reader had 
splits").isNotEmpty();
+
+            enumerator.addSplitsBack(splitsOfFailingReader, failingReader);
+
+            for (int reader = 0; reader < failingReader; reader++) {
+                assertThat(context.getSignalCount(reader))
+                        .as("reader %s must not be signalled again after it 
was already told", reader)

Review Comment:
   `spotless:check` fails on this line, CI will stop here.
   
   ```suggestion
                           .as(
                                   "reader %s must not be signalled again after 
it was already told",
                                   reader)
   ```



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