From 7a353a1c27323dc74319234a9285e3673b3d9616 Mon Sep 17 00:00:00 2001 From: Jennifer Xiong Date: Fri, 20 Mar 2026 13:09:25 -0700 Subject: [PATCH 1/3] new fixes --- .../status/FlinkBlueGreenDeploymentStatus.java | 6 ++++++ .../bluegreen/BlueGreenDeploymentService.java | 5 +++++ .../operator/utils/bluegreen/BlueGreenUtils.java | 16 ++++++++++++++++ .../FlinkBlueGreenDeploymentControllerTest.java | 7 +++++++ 4 files changed, 34 insertions(+) diff --git a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java index 104f195efa..2dcfc3eb55 100644 --- a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java +++ b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java @@ -44,6 +44,12 @@ public class FlinkBlueGreenDeploymentStatus { /** Last reconciled (serialized) deployment spec. */ private String lastReconciledSpec; + /** + * Previous reconciled spec, saved before a transition overwrites lastReconciledSpec. Used to + * restore lastReconciledSpec on abort so it stays consistent with the active child. + */ + private String previousReconciledSpec; + /** Timestamp of last reconciliation. */ private String lastReconciledTimestamp; 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..9355f09a2c 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 @@ -58,6 +58,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.revertLastReconciledSpec; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.setLastReconciledSpec; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.triggerSavepoint; @@ -697,6 +698,10 @@ private UpdateControl abortDeployment( context.getDeploymentStatus().setBlueGreenState(previousState); context.getDeploymentStatus().setSavepointTriggerId(null); + // Revert lastReconciledSpec to the pre-transition spec so it stays + // consistent with the active child that is still running + revertLastReconciledSpec(context); + var error = String.format( "Aborting deployment '%s', rolling B/G deployment back to %s", diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java index af3ccbd65c..4cdd169a4f 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java @@ -89,6 +89,8 @@ public static BlueGreenDiffType getSpecDiff(BlueGreenContext context) { public static void setLastReconciledSpec(BlueGreenContext context) { FlinkBlueGreenDeploymentStatus deploymentStatus = context.getDeploymentStatus(); + // Save the current lastReconciledSpec so it can be restored on abort + deploymentStatus.setPreviousReconciledSpec(deploymentStatus.getLastReconciledSpec()); deploymentStatus.setLastReconciledSpec( SpecUtils.writeSpecAsJSON(context.getBgDeployment().getSpec(), "spec")); deploymentStatus.setLastReconciledTimestamp(Instant.now().toString()); @@ -104,6 +106,20 @@ public static void revertToLastSpec(BlueGreenContext context) { replaceFlinkBlueGreenDeployment(context); } + /** + * Restores lastReconciledSpec from previousReconciledSpec and reverts the B/G CR's live spec to + * match. Called on abort so lastReconciledSpec stays consistent with the active child. + */ + public static void revertLastReconciledSpec(BlueGreenContext context) { + FlinkBlueGreenDeploymentStatus deploymentStatus = context.getDeploymentStatus(); + String previousSpec = deploymentStatus.getPreviousReconciledSpec(); + if (previousSpec != null) { + deploymentStatus.setLastReconciledSpec(previousSpec); + deploymentStatus.setPreviousReconciledSpec(null); + revertToLastSpec(context); + } + } + /** * Extracts a configuration option value from the Blue/Green deployment spec. * 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..b6d09c0a93 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,13 @@ 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"); + assertNull(rs.reconciledStatus.getPreviousReconciledSpec()); // Simulate another change in the spec to trigger a redeployment customValue = UUID.randomUUID().toString(); From 2bf72e2a8d931957b3f487ddabff169133b31bd7 Mon Sep 17 00:00:00 2001 From: Jennifer Xiong Date: Fri, 20 Mar 2026 13:09:25 -0700 Subject: [PATCH 2/3] new fixes --- .../controller/bluegreen/BlueGreenDeploymentService.java | 4 ---- 1 file changed, 4 deletions(-) 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 9355f09a2c..48e08ff56d 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 @@ -169,10 +169,6 @@ public UpdateControl checkAndInitiateDeployment( 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) { From 7c594f92db508c7e0f7773e3c7475aee38d06cde Mon Sep 17 00:00:00 2001 From: Jennifer Xiong Date: Fri, 20 Mar 2026 13:09:25 -0700 Subject: [PATCH 3/3] new fixes --- .../FlinkBlueGreenDeploymentStatus.java | 6 ---- .../bluegreen/BlueGreenDeploymentService.java | 35 ++++++++++++++++--- .../utils/bluegreen/BlueGreenUtils.java | 16 --------- ...linkBlueGreenDeploymentControllerTest.java | 1 - 4 files changed, 31 insertions(+), 27 deletions(-) diff --git a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java index 2dcfc3eb55..104f195efa 100644 --- a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java +++ b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkBlueGreenDeploymentStatus.java @@ -44,12 +44,6 @@ public class FlinkBlueGreenDeploymentStatus { /** Last reconciled (serialized) deployment spec. */ private String lastReconciledSpec; - /** - * Previous reconciled spec, saved before a transition overwrites lastReconciledSpec. Used to - * restore lastReconciledSpec on abort so it stays consistent with the active child. - */ - private String previousReconciledSpec; - /** Timestamp of last reconciliation. */ private String lastReconciledTimestamp; 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 48e08ff56d..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,7 +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.revertLastReconciledSpec; +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; @@ -68,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 ==================== /** @@ -164,6 +175,8 @@ public UpdateControl checkAndInitiateDeployment( } try { + // Save pre-transition lastReconciledSpec so we can restore on abort + savePreTransitionSpec(context); var result = startTransition( context, @@ -186,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); @@ -694,9 +709,14 @@ private UpdateControl abortDeployment( context.getDeploymentStatus().setBlueGreenState(previousState); context.getDeploymentStatus().setSavepointTriggerId(null); - // Revert lastReconciledSpec to the pre-transition spec so it stays - // consistent with the active child that is still running - revertLastReconciledSpec(context); + // 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( @@ -705,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) { @@ -745,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/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java index 4cdd169a4f..af3ccbd65c 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtils.java @@ -89,8 +89,6 @@ public static BlueGreenDiffType getSpecDiff(BlueGreenContext context) { public static void setLastReconciledSpec(BlueGreenContext context) { FlinkBlueGreenDeploymentStatus deploymentStatus = context.getDeploymentStatus(); - // Save the current lastReconciledSpec so it can be restored on abort - deploymentStatus.setPreviousReconciledSpec(deploymentStatus.getLastReconciledSpec()); deploymentStatus.setLastReconciledSpec( SpecUtils.writeSpecAsJSON(context.getBgDeployment().getSpec(), "spec")); deploymentStatus.setLastReconciledTimestamp(Instant.now().toString()); @@ -106,20 +104,6 @@ public static void revertToLastSpec(BlueGreenContext context) { replaceFlinkBlueGreenDeployment(context); } - /** - * Restores lastReconciledSpec from previousReconciledSpec and reverts the B/G CR's live spec to - * match. Called on abort so lastReconciledSpec stays consistent with the active child. - */ - public static void revertLastReconciledSpec(BlueGreenContext context) { - FlinkBlueGreenDeploymentStatus deploymentStatus = context.getDeploymentStatus(); - String previousSpec = deploymentStatus.getPreviousReconciledSpec(); - if (previousSpec != null) { - deploymentStatus.setLastReconciledSpec(previousSpec); - deploymentStatus.setPreviousReconciledSpec(null); - revertToLastSpec(context); - } - } - /** * Extracts a configuration option value from the Blue/Green deployment spec. * 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 b6d09c0a93..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 @@ -555,7 +555,6 @@ public void verifyFailureDuringTransition(FlinkVersion flinkVersion) throws Exce assertFalse( rs.reconciledStatus.getLastReconciledSpec().contains(customValue), "lastReconciledSpec should be reverted on abort to match the active child"); - assertNull(rs.reconciledStatus.getPreviousReconciledSpec()); // Simulate another change in the spec to trigger a redeployment customValue = UUID.randomUUID().toString();