From b6b93a63c49c7dbf56e817fcbbd7bfbd4bc0eec0 Mon Sep 17 00:00:00 2001 From: Li Guo Date: Tue, 1 Sep 2026 11:23:28 +0000 Subject: [PATCH] [FLINK-39108][checkpoint] Respect disabled interval-during-backlog for the first checkpoint When a source reported isProcessingBacklog=true while execution.checkpointing.interval-during-backlog was disabled, the periodic trigger armed by startCheckpointScheduler stayed scheduled: setIsProcessingBacklog only reschedules when the new effective interval is shorter and skips the disabled case entirely. On top of that, ScheduledTrigger#run stopped future scheduling when the effective interval was disabled but still triggered the current run. As a result at least one checkpoint fired during backlog processing even though checkpointing should have been suspended. Cancel the armed periodic trigger when the effective interval becomes disabled, and make ScheduledTrigger re-check the effective interval before triggering so a pending run is skipped once checkpointing is disabled. The trigger is re-armed by setIsProcessingBacklog when the backlog ends. This also re-enables CheckpointIntervalDuringBacklogITCase#testNoCheckpointDuringBacklog, which was disabled waiting for this fix (its annotation referenced FLINK-39018, a typo for FLINK-39108). --- .../checkpoint/CheckpointCoordinator.java | 11 ++ .../CheckpointCoordinatorTriggeringTest.java | 135 ++++++++++++++++++ ...CheckpointIntervalDuringBacklogITCase.java | 2 - 3 files changed, 146 insertions(+), 2 deletions(-) diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java b/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java index 7750543be41001..73f9375c4ebe55 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinator.java @@ -483,6 +483,12 @@ public void setIsProcessingBacklog(OperatorID operatorID, boolean isProcessingBa < nextCheckpointTriggeringRelativeTime) { rescheduleTrigger(currentRelativeTime, currentCheckpointInterval); } + } else { + // Periodic checkpointing is disabled for the current backlog status. Cancel any + // trigger that has already been scheduled (e.g. by the checkpoint scheduler before + // the backlog was reported) so that no checkpoint is triggered until the status + // changes again (FLINK-39108). + cancelPeriodicTrigger(); } } } @@ -2224,9 +2230,14 @@ public void run() { - clock.relativeTimeMillis()), TimeUnit.MILLISECONDS); } else { + // Periodic checkpointing has been disabled in the meantime (e.g. a source + // reported backlog with a disabled interval-during-backlog). Do not trigger + // this run either; the trigger will be rescheduled when the effective + // interval becomes enabled again (FLINK-39108). nextCheckpointTriggeringRelativeTime = Long.MAX_VALUE; currentPeriodicTrigger = null; currentPeriodicTriggerFuture = null; + return; } } diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinatorTriggeringTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinatorTriggeringTest.java index 12b404bf42a7d0..c68a8d9fcd2b63 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinatorTriggeringTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/CheckpointCoordinatorTriggeringTest.java @@ -29,6 +29,7 @@ import org.apache.flink.runtime.executiongraph.ExecutionGraph; import org.apache.flink.runtime.executiongraph.ExecutionVertex; import org.apache.flink.runtime.jobgraph.JobVertexID; +import org.apache.flink.runtime.jobgraph.OperatorID; import org.apache.flink.runtime.jobgraph.SavepointRestoreSettings; import org.apache.flink.runtime.jobgraph.tasks.CheckpointCoordinatorConfiguration; import org.apache.flink.runtime.jobgraph.tasks.CheckpointCoordinatorConfiguration.CheckpointCoordinatorConfigurationBuilder; @@ -39,6 +40,7 @@ import org.apache.flink.testutils.junit.utils.TempDirUtils; import org.apache.flink.util.ExceptionUtils; import org.apache.flink.util.ExecutorUtils; +import org.apache.flink.util.clock.ManualClock; import org.apache.flink.util.concurrent.FutureUtils; import org.apache.flink.util.concurrent.ManuallyTriggeredScheduledExecutor; import org.apache.flink.util.concurrent.ScheduledExecutorServiceAdapter; @@ -184,6 +186,139 @@ private void checkRecordedTriggeredCheckpoints( } } + /** + * Tests that no periodic checkpoint is triggered while a source is processing backlog and + * {@code execution.checkpointing.interval-during-backlog} is disabled, even if the first + * periodic trigger was already scheduled before the source reported the backlog (FLINK-39108). + */ + @Test + void testFirstScheduledCheckpointNotTriggeredWhenBacklogCheckpointingDisabled() + throws Exception { + CheckpointCoordinatorTestingUtils.CheckpointRecorderTaskManagerGateway gateway = + new CheckpointCoordinatorTestingUtils.CheckpointRecorderTaskManagerGateway(); + + JobVertexID jobVertexID = new JobVertexID(); + ExecutionGraph graph = + new CheckpointCoordinatorTestingUtils.CheckpointExecutionGraphBuilder() + .addJobVertex(jobVertexID) + .setTaskManagerGateway(gateway) + .build(EXECUTOR_RESOURCE.getExecutor()); + + ExecutionVertex vertex = graph.getJobVertex(jobVertexID).getTaskVertices()[0]; + ExecutionAttemptID attemptID = vertex.getCurrentExecutionAttempt().getAttemptId(); + + CheckpointCoordinatorConfiguration checkpointCoordinatorConfiguration = + new CheckpointCoordinatorConfigurationBuilder() + .setCheckpointInterval(10) + .setCheckpointIntervalDuringBacklog( + CheckpointCoordinatorConfiguration.DISABLED_CHECKPOINT_INTERVAL) + .setCheckpointTimeout(200000) + .setMaxConcurrentCheckpoints(Integer.MAX_VALUE) + .build(); + CheckpointCoordinator checkpointCoordinator = + new CheckpointCoordinatorBuilder() + .setCheckpointCoordinatorConfiguration(checkpointCoordinatorConfiguration) + .setTimer(manuallyTriggeredScheduledExecutor) + .setClock(new ManualClock()) + .build(graph); + + try { + checkpointCoordinator.startCheckpointScheduler(); + + // a source reports isProcessingBacklog=true before the first periodic checkpoint + // has been triggered; with a disabled backlog interval no checkpoint may be + // triggered until the backlog is processed + OperatorID operatorID = new OperatorID(); + checkpointCoordinator.setIsProcessingBacklog(operatorID, true); + assertThat(checkpointCoordinator.isCurrentPeriodicTriggerAvailable()) + .as("the armed periodic trigger is cancelled right away") + .isFalse(); + + manuallyTriggeredScheduledExecutor.triggerNonPeriodicScheduledTasks( + CheckpointCoordinator.ScheduledTrigger.class); + manuallyTriggeredScheduledExecutor.triggerAll(); + + assertThat(checkpointCoordinator.getNumberOfPendingCheckpoints()) + .as( + "No checkpoint should be triggered during backlog if the checkpoint " + + "interval during backlog is disabled") + .isZero(); + assertThat(gateway.getTriggeredCheckpoints(attemptID)).isEmpty(); + + // the backlog ends: the periodic trigger is re-armed and checkpointing resumes + checkpointCoordinator.setIsProcessingBacklog(operatorID, false); + assertThat(checkpointCoordinator.isCurrentPeriodicTriggerAvailable()).isTrue(); + manuallyTriggeredScheduledExecutor.triggerNonPeriodicScheduledTasks( + CheckpointCoordinator.ScheduledTrigger.class); + manuallyTriggeredScheduledExecutor.triggerAll(); + assertThat(gateway.getTriggeredCheckpoints(attemptID)).hasSize(1); + } finally { + checkpointCoordinator.shutdown(); + } + } + + /** + * Tests that a periodic trigger armed while a source is already processing backlog does not + * trigger a checkpoint if {@code execution.checkpointing.interval-during-backlog} is disabled, + * and that checkpointing resumes once the backlog is processed (FLINK-39108). + */ + @Test + void testSchedulerStartedDuringBacklogDoesNotTriggerWhenBacklogCheckpointingDisabled() + throws Exception { + CheckpointCoordinatorTestingUtils.CheckpointRecorderTaskManagerGateway gateway = + new CheckpointCoordinatorTestingUtils.CheckpointRecorderTaskManagerGateway(); + + JobVertexID jobVertexID = new JobVertexID(); + ExecutionGraph graph = + new CheckpointCoordinatorTestingUtils.CheckpointExecutionGraphBuilder() + .addJobVertex(jobVertexID) + .setTaskManagerGateway(gateway) + .build(EXECUTOR_RESOURCE.getExecutor()); + + ExecutionVertex vertex = graph.getJobVertex(jobVertexID).getTaskVertices()[0]; + ExecutionAttemptID attemptID = vertex.getCurrentExecutionAttempt().getAttemptId(); + + CheckpointCoordinatorConfiguration checkpointCoordinatorConfiguration = + new CheckpointCoordinatorConfigurationBuilder() + .setCheckpointInterval(10) + .setCheckpointIntervalDuringBacklog( + CheckpointCoordinatorConfiguration.DISABLED_CHECKPOINT_INTERVAL) + .setCheckpointTimeout(200000) + .setMaxConcurrentCheckpoints(Integer.MAX_VALUE) + .build(); + CheckpointCoordinator checkpointCoordinator = + new CheckpointCoordinatorBuilder() + .setCheckpointCoordinatorConfiguration(checkpointCoordinatorConfiguration) + .setTimer(manuallyTriggeredScheduledExecutor) + .setClock(new ManualClock()) + .build(graph); + + try { + // the source reports backlog first, the scheduler is started afterwards (this is the + // order in a running job: the enumerator reports on start, the scheduler starts once + // all tasks are running) + OperatorID operatorID = new OperatorID(); + checkpointCoordinator.setIsProcessingBacklog(operatorID, true); + checkpointCoordinator.startCheckpointScheduler(); + + manuallyTriggeredScheduledExecutor.triggerNonPeriodicScheduledTasks( + CheckpointCoordinator.ScheduledTrigger.class); + manuallyTriggeredScheduledExecutor.triggerAll(); + + assertThat(checkpointCoordinator.getNumberOfPendingCheckpoints()).isZero(); + assertThat(gateway.getTriggeredCheckpoints(attemptID)).isEmpty(); + assertThat(checkpointCoordinator.isCurrentPeriodicTriggerAvailable()).isFalse(); + + checkpointCoordinator.setIsProcessingBacklog(operatorID, false); + manuallyTriggeredScheduledExecutor.triggerNonPeriodicScheduledTasks( + CheckpointCoordinator.ScheduledTrigger.class); + manuallyTriggeredScheduledExecutor.triggerAll(); + assertThat(gateway.getTriggeredCheckpoints(attemptID)).hasSize(1); + } finally { + checkpointCoordinator.shutdown(); + } + } + @Test void testTriggeringFullSnapshotAfterJobmasterFailover() throws Exception { CheckpointCoordinatorTestingUtils.CheckpointRecorderTaskManagerGateway gateway = diff --git a/flink-tests/src/test/java/org/apache/flink/test/checkpointing/CheckpointIntervalDuringBacklogITCase.java b/flink-tests/src/test/java/org/apache/flink/test/checkpointing/CheckpointIntervalDuringBacklogITCase.java index c3ebd46f304e07..fb7950ecd5a787 100644 --- a/flink-tests/src/test/java/org/apache/flink/test/checkpointing/CheckpointIntervalDuringBacklogITCase.java +++ b/flink-tests/src/test/java/org/apache/flink/test/checkpointing/CheckpointIntervalDuringBacklogITCase.java @@ -44,7 +44,6 @@ import org.apache.flink.util.CloseableIterator; import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; @@ -140,7 +139,6 @@ void testDefaultCheckpointIntervalDuringBacklog() throws Exception { } @Test - @Disabled("FLINK-39018") // FLINK-39108 void testNoCheckpointDuringBacklog() throws Exception { final int recordsBeforeSwitch = NUM_RECORDS / 2; Duration expectedSwitchTime = Duration.ofMillis(recordsBeforeSwitch * SLEEP_MS_PER_RECORD);