SteveStevenpoor commented on code in PR #28689:
URL: https://github.com/apache/flink/pull/28689#discussion_r3593630420
##########
flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/SourceOperatorSplitWatermarkAlignmentTest.java:
##########
@@ -569,6 +569,58 @@ void testAlignmentCheckIsDeferredForIdleSplits() throws
Exception {
0L,
operator.getSplitMetricGroup(split0.splitId()).getAccumulatedPausedTime());
}
+ @Test
+ void testPausedIdleSplitsCanBeResumedByAlignmentCheck() throws Exception {
+ final long idleTimeout = 100;
+ final MockSourceReader sourceReader =
+ new MockSourceReader(WaitingForSplits.DO_NOT_WAIT_FOR_SPLITS,
true, true);
+ final TestProcessingTimeService processingTimeService = new
TestProcessingTimeService();
+ final SourceOperator<Integer, MockSourceSplit> operator =
+ createAndOpenSourceOperatorWithIdleness(
+ sourceReader, processingTimeService, idleTimeout);
+
+ final MockSourceSplit split0 = new MockSourceSplit(0, 0, 10);
+ final int allowedWatermark4 = 4;
+ final int allowedWatermark7 = 7;
+ split0.addRecord(4);
+ split0.addRecord(5);
+ split0.addRecord(6);
+ split0.addRecord(7);
+ split0.addRecord(8);
+ operator.handleOperatorEvent(
+ new AddSplitEvent<>(Arrays.asList(split0), new
MockSourceSplitSerializer()));
+ final CollectingDataOutput<Integer> actualOutput = new
CollectingDataOutput<>();
+
+ // Emit enough records to fill the sampler buffer
+ for (int i = 0; i < WATERMARK_ALIGNMENT_BUFFER_SIZE.defaultValue();
i++) {
+ operator.emitNext(actualOutput);
+ processingTimeService.advance(idleTimeout - 1);
+ }
+ sampleAllWatermarks(processingTimeService);
+ assertOutput(actualOutput, Arrays.asList(4, 5, 6));
+
+ // Alignment check fires and pauses the split
+ operator.handleOperatorEvent(new
WatermarkAlignmentEvent(allowedWatermark4));
+
assertThat(operator.getSplitMetricGroup(split0.splitId()).isPaused()).isTrue();
+ assertThat(sourceReader.getPausedSplits()).containsExactly("0");
+ assertOutput(actualOutput, Arrays.asList(4, 5, 6));
+
+ // Normally idlenessTimer can't elapse while the split is paused
+ // So calling it manually to simulate a race condition
+ operator.updateCurrentSplitIdle(split0.splitId(), true);
Review Comment:
> But the test as it is also provides cover for this bug, doesn't it?
Yeah, the current tests covers the problematic final state. My point is only
that the simulated behavior is slightly different from the prod scenario.
The suggested injection would make the test a bit more precise and focused
on the race window described in the issue. That said, I don't think its
strictly necessary if we only want to cover the resulting state.
--
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]