This is an automated email from the ASF dual-hosted git repository.
1996fanrui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/master by this push:
new 8ce71ee0b9c [FLINK-40518][checkpointing] Only enable
checkpointing-during-recovery when unaligned checkpoints are enabled
8ce71ee0b9c is described below
commit 8ce71ee0b9ce3dcc53711326a2a92338dfe9c185
Author: Rui Fan <[email protected]>
AuthorDate: Mon Aug 31 17:28:20 2026 +0200
[FLINK-40518][checkpointing] Only enable checkpointing-during-recovery when
unaligned checkpoints are enabled
---
.../flink/configuration/CheckpointingOptions.java | 10 +++----
.../configuration/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 41397521f84..b36f9767ed4 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 class CheckpointingOptions {
/**
* Determines whether unaligned checkpoint support during recovery is
enabled.
*
- * <p>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.
+ * <p>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 class CheckpointingOptions {
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 2334c71fd29..9c8940eb246 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 @@ class CheckpointingOptionsTest {
.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(