diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java index 6274881df2..ad6f0d9bfa 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenDeploymentService.java @@ -44,7 +44,9 @@ import org.slf4j.LoggerFactory; import java.time.Instant; +import java.util.Map; import java.util.Objects; +import java.util.concurrent.ConcurrentHashMap; import static org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenKubernetesService.deleteFlinkDeployment; import static org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenKubernetesService.deployCluster; @@ -58,6 +60,7 @@ import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.isSavepointRequired; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.millisToInstantStr; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.prepareFlinkDeployment; +import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.revertToLastSpec; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.setLastReconciledSpec; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.triggerSavepoint; @@ -67,6 +70,15 @@ public class BlueGreenDeploymentService { private static final Logger LOG = LoggerFactory.getLogger(BlueGreenDeploymentService.class); private static final long RETRY_DELAY_MS = 500; + /** + * In-memory cache of the pre-transition lastReconciledSpec, keyed by namespace. Saved before + * setLastReconciledSpec overwrites it during a transition, restored on abort so + * lastReconciledSpec stays consistent with the active child. Lost on operator restart, which + * falls back to pre-fix behavior (no revert) — acceptable since restart during the abort window + * is rare and not a regression. + */ + private final Map previousReconciledSpecs = new ConcurrentHashMap<>(); + // ==================== Deployment Initiation Methods ==================== /** @@ -163,15 +175,13 @@ public UpdateControl checkAndInitiateDeployment( } try { + // Save pre-transition lastReconciledSpec so we can restore on abort + savePreTransitionSpec(context); var result = startTransition( context, currentBlueGreenDeploymentType, currentFlinkDeployment); - // Only stamp lastReconciledSpec after the transition - // succeeds. If stamped before and the transition - // fails/aborts, lastReconciledSpec drifts from the - // active child's actual spec setLastReconciledSpec(context); return result; } catch (Exception e) { @@ -189,6 +199,8 @@ public UpdateControl checkAndInitiateDeployment( context.getBgDeployment().getMetadata().getName(), Objects.toString(jobSpec.getInitialSavepointPath(), "")); try { + // Save pre-transition lastReconciledSpec so we can restore on abort + savePreTransitionSpec(context); var result = startSavepointRedeployTransition( context, currentBlueGreenDeploymentType); @@ -697,6 +709,15 @@ private UpdateControl abortDeployment( context.getDeploymentStatus().setBlueGreenState(previousState); context.getDeploymentStatus().setSavepointTriggerId(null); + // Restore lastReconciledSpec and the B/G CR spec to the pre-transition state + // so they stay consistent with the active child that is still running. + String namespace = context.getBgDeployment().getMetadata().getNamespace(); + String previousSpec = previousReconciledSpecs.remove(namespace); + if (previousSpec != null) { + context.getDeploymentStatus().setLastReconciledSpec(previousSpec); + revertToLastSpec(context); + } + var error = String.format( "Aborting deployment '%s', rolling B/G deployment back to %s", @@ -704,6 +725,12 @@ private UpdateControl abortDeployment( return markDeploymentFailing(context, error); } + private void savePreTransitionSpec(BlueGreenContext context) { + String namespace = context.getBgDeployment().getMetadata().getNamespace(); + previousReconciledSpecs.put( + namespace, context.getDeploymentStatus().getLastReconciledSpec()); + } + @NotNull private static UpdateControl markDeploymentFailing( BlueGreenContext context, String error) { @@ -744,6 +771,7 @@ public UpdateControl finalizeBlueGreenDeployment( context.getDeploymentStatus().setDeploymentReadyTimestamp(millisToInstantStr(0)); context.getDeploymentStatus().setAbortTimestamp(millisToInstantStr(0)); context.getDeploymentStatus().setSavepointTriggerId(null); + previousReconciledSpecs.remove(context.getBgDeployment().getMetadata().getNamespace()); updateBlueGreenIngress(context, nextState); diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java index e657bb7227..c18fafa24b 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkBlueGreenDeploymentControllerTest.java @@ -549,6 +549,12 @@ public void verifyFailureDuringTransition(FlinkVersion flinkVersion) throws Exce // savepointTriggerId must be cleared on abort so the next transition // triggers a fresh savepoint instead of reusing a stale triggerId assertNull(rs.reconciledStatus.getSavepointTriggerId()); + // lastReconciledSpec must NOT contain the failed spec change — on abort + // it should be reverted to the pre-transition spec so it stays consistent + // with the active child that is still running + assertFalse( + rs.reconciledStatus.getLastReconciledSpec().contains(customValue), + "lastReconciledSpec should be reverted on abort to match the active child"); // Simulate another change in the spec to trigger a redeployment customValue = UUID.randomUUID().toString();