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]