Efrat19 commented on code in PR #28689:
URL: https://github.com/apache/flink/pull/28689#discussion_r3585574816
##########
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
Review Comment:
I updated the test to reproduce the race.
On master version the race still reproduces, but the test fails expecting
`sourceReader.getPausedSplits()` to be empty (which is what this fix is about)
```
[ERROR] Failures:
[ERROR]
SourceOperatorSplitWatermarkAlignmentTest.testPausedIdleSplitsCanBeResumedByAlignmentCheck:624
Expecting empty but was: ["0"]
[INFO]
[ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0
[INFO]
```
--
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]