rkhachatryan commented on a change in pull request #13040:
URL: https://github.com/apache/flink/pull/13040#discussion_r495781246
##########
File path:
flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/ZooKeeperCompletedCheckpointStoreITCase.java
##########
@@ -283,6 +286,54 @@ public void testConcurrentCheckpointOperations() throws
Exception {
recoveredTestCheckpoint.awaitDiscard();
}
+ /**
+ * FLINK-17073 tests that there is no request triggered when there are
too many checkpoints
+ * waiting to clean and that it resumes when the number of waiting
checkpoints as gone below
+ * the threshold.
+ *
+ */
+ @Test
+ public void testChekpointingPausesAndResumeWhenTooManyCheckpoints()
throws Exception{
+ ManualClock clock = new ManualClock();
+ clock.advanceTime(1, TimeUnit.DAYS);
+ int maxCleaningCheckpoints = 1;
+ CheckpointsCleaner checkpointsCleaner = new
CheckpointsCleaner();
+ CheckpointRequestDecider checkpointRequestDecider = new
CheckpointRequestDecider(maxCleaningCheckpoints, unused ->{}, clock, 1, new
AtomicInteger(0)::get, checkpointsCleaner::getNumberOfCheckpointsToClean);
+
+ final int maxCheckpointsToRetain = 1;
+ Executors.PausableThreadPoolExecutor executor =
Executors.pausableExecutor();
+ ZooKeeperCompletedCheckpointStore checkpointStore =
createCompletedCheckpoints(maxCheckpointsToRetain, executor);
+
+ //pause the executor to pause checkpoints cleaning, to allow
assertions
+ executor.pause();
+
+ int nbCheckpointsToInject = 3;
+ for (int i = 1; i <= nbCheckpointsToInject; i++) {
+ // add checkpoints to clean
+ TestCompletedCheckpoint completedCheckpoint = new
TestCompletedCheckpoint(new JobID(), i,
+ i, Collections.emptyMap(),
CheckpointProperties.forCheckpoint(CheckpointRetentionPolicy.RETAIN_ON_FAILURE),
+ checkpointsCleaner::cleanCheckpoint);
+ checkpointStore.addCheckpoint(completedCheckpoint);
+ }
+
+ Thread.sleep(100L); // give time to submit checkpoints for
cleaning
+
+ int nbCheckpointsSubmittedForCleaningByCheckpointStore =
nbCheckpointsToInject - maxCheckpointsToRetain;
+
assertEquals(nbCheckpointsSubmittedForCleaningByCheckpointStore,
checkpointsCleaner.getNumberOfCheckpointsToClean());
Review comment:
I think forcing the developer to change the tests should be avoided if
possible, and here it is possible.
Also, condition on the public interface rather than internal details makes
expectations more clear.
From my side, the PR can be merged after resolving this issue and grooming
commits.
Probably some other developers would like to take a brief look - I'll ask
about it.
----------------------------------------------------------------
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.
For queries about this service, please contact Infrastructure at:
[email protected]