Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -831,10 +831,9 @@ public static boolean isUnalignedCheckpointInterruptibleTimersEnabled(Configurat
/**
* 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}
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down