From acf10683a5b16581f4b7a551aa1c23afce45ca48 Mon Sep 17 00:00:00 2001 From: Rui Fan <1996fanrui@gmail.com> Date: Mon, 31 Aug 2026 17:28:20 +0200 Subject: [PATCH] [FLINK-40518][checkpointing] Only enable checkpointing-during-recovery when unaligned checkpoints are enabled --- .../configuration/CheckpointingOptions.java | 10 +++--- .../CheckpointingOptionsTest.java | 31 +++++++++++++++---- 2 files changed, 30 insertions(+), 11 deletions(-) diff --git a/flink-core/src/main/java/org/apache/flink/configuration/CheckpointingOptions.java b/flink-core/src/main/java/org/apache/flink/configuration/CheckpointingOptions.java index 41397521f8452f..b36f9767ed4c7f 100644 --- a/flink-core/src/main/java/org/apache/flink/configuration/CheckpointingOptions.java +++ b/flink-core/src/main/java/org/apache/flink/configuration/CheckpointingOptions.java @@ -831,10 +831,9 @@ public static boolean isUnalignedCheckpointInterruptibleTimersEnabled(Configurat /** * Determines whether unaligned checkpoint support during recovery is enabled. * - *

This feature requires {@link #UNALIGNED_RECOVER_OUTPUT_ON_DOWNSTREAM} to be enabled. Note - * that it does not require unaligned checkpoints to be currently enabled, because a job may - * restore from an unaligned checkpoint while having unaligned checkpoints disabled for the new - * execution. + *

Requires both {@link #UNALIGNED_RECOVER_OUTPUT_ON_DOWNSTREAM} and unaligned checkpoints to + * be enabled, because checkpointing during recovery is only supported on the unaligned + * barrier-handler path. * * @param config the configuration to check * @return {@code true} if unaligned checkpointing during recovery is enabled, {@code false} @@ -845,6 +844,7 @@ public static boolean isCheckpointingDuringRecoveryEnabled(Configuration config) if (!config.get(UNALIGNED_RECOVER_OUTPUT_ON_DOWNSTREAM)) { return false; } - return config.get(CHECKPOINTING_DURING_RECOVERY_ENABLED); + return config.get(CHECKPOINTING_DURING_RECOVERY_ENABLED) + && isUnalignedCheckpointEnabled(config); } } diff --git a/flink-core/src/test/java/org/apache/flink/configuration/CheckpointingOptionsTest.java b/flink-core/src/test/java/org/apache/flink/configuration/CheckpointingOptionsTest.java index 2334c71fd2922c..9c8940eb2462d6 100644 --- a/flink-core/src/test/java/org/apache/flink/configuration/CheckpointingOptionsTest.java +++ b/flink-core/src/test/java/org/apache/flink/configuration/CheckpointingOptionsTest.java @@ -358,15 +358,34 @@ void testIsCheckpointingDuringRecoveryEnabled() { .as("During-recovery should be disabled when during-recovery option is not enabled") .isFalse(); - // Test when both options are enabled - should return true - Configuration bothEnabledConfig = new Configuration(); - bothEnabledConfig.set(CheckpointingOptions.UNALIGNED_RECOVER_OUTPUT_ON_DOWNSTREAM, true); - bothEnabledConfig.set(CheckpointingOptions.CHECKPOINTING_DURING_RECOVERY_ENABLED, true); - assertThat(CheckpointingOptions.isCheckpointingDuringRecoveryEnabled(bothEnabledConfig)) + // Test when all three prerequisites are enabled - should return true + Configuration allEnabledConfig = new Configuration(); + allEnabledConfig.set(CheckpointingOptions.CHECKPOINTING_INTERVAL, Duration.ofSeconds(5)); + allEnabledConfig.set(CheckpointingOptions.UNALIGNED_RECOVER_OUTPUT_ON_DOWNSTREAM, true); + allEnabledConfig.set(CheckpointingOptions.CHECKPOINTING_DURING_RECOVERY_ENABLED, true); + allEnabledConfig.set(CheckpointingOptions.ENABLE_UNALIGNED, true); + assertThat(CheckpointingOptions.isCheckpointingDuringRecoveryEnabled(allEnabledConfig)) .as( - "During-recovery should be enabled when both recover-output-on-downstream and during-recovery are enabled") + "During-recovery should be enabled when recover-output-on-downstream, during-recovery and unaligned checkpoints are all enabled") .isTrue(); + // Test when recover-output-on-downstream and during-recovery are enabled but unaligned + // checkpoints are disabled - should return false (checkpointing during recovery only works + // on the unaligned barrier-handler path). + Configuration unalignedDisabledConfig = new Configuration(); + unalignedDisabledConfig.set( + CheckpointingOptions.CHECKPOINTING_INTERVAL, Duration.ofSeconds(5)); + unalignedDisabledConfig.set( + CheckpointingOptions.UNALIGNED_RECOVER_OUTPUT_ON_DOWNSTREAM, true); + unalignedDisabledConfig.set( + CheckpointingOptions.CHECKPOINTING_DURING_RECOVERY_ENABLED, true); + unalignedDisabledConfig.set(CheckpointingOptions.ENABLE_UNALIGNED, false); + assertThat( + CheckpointingOptions.isCheckpointingDuringRecoveryEnabled( + unalignedDisabledConfig)) + .as("During-recovery should be disabled when unaligned checkpoints are disabled") + .isFalse(); + // Test when recover-output-on-downstream is explicitly false and during-recovery is true Configuration explicitlyDisabledConfig = new Configuration(); explicitlyDisabledConfig.set(