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 9543fecabd..560f27473c 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 @@ -22,14 +22,21 @@ import org.apache.flink.kubernetes.operator.api.FlinkDeployment; import org.apache.flink.kubernetes.operator.api.bluegreen.BlueGreenDeploymentType; import org.apache.flink.kubernetes.operator.api.bluegreen.BlueGreenDiffType; +import org.apache.flink.kubernetes.operator.api.diff.DiffType; import org.apache.flink.kubernetes.operator.api.lifecycle.ResourceLifecycleState; +import org.apache.flink.kubernetes.operator.api.spec.FlinkBlueGreenDeploymentSpec; +import org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec; +import org.apache.flink.kubernetes.operator.api.spec.JobState; +import org.apache.flink.kubernetes.operator.api.spec.KubernetesDeploymentMode; import org.apache.flink.kubernetes.operator.api.status.FlinkBlueGreenDeploymentState; import org.apache.flink.kubernetes.operator.api.status.Savepoint; import org.apache.flink.kubernetes.operator.api.status.SavepointFormatType; import org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType; +import org.apache.flink.kubernetes.operator.api.utils.SpecUtils; import org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions; import org.apache.flink.kubernetes.operator.controller.FlinkBlueGreenDeployments; import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext; +import org.apache.flink.kubernetes.operator.reconciler.diff.ReflectiveDiffBuilder; import org.apache.flink.kubernetes.operator.utils.IngressUtils; import org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils; import org.apache.flink.util.Preconditions; @@ -53,7 +60,6 @@ import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.fetchSavepointInfo; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.getReconciliationReschedInterval; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.getSpecDiff; -import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.hasSpecChanged; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.instantStrToMillis; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.isSavepointRequired; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.millisToInstantStr; @@ -69,22 +75,33 @@ public class BlueGreenDeploymentService { // ==================== Deployment Initiation Methods ==================== - /** - * Initiates a new Blue/Green deployment. - * - * @param context the transition context - * @param nextBlueGreenDeploymentType the type of deployment to create - * @param nextState the next state to transition to - * @param lastCheckpoint the checkpoint to restore from (can be null) - * @param isFirstDeployment whether this is the first deployment - * @return UpdateControl for the deployment - */ + /** Convenience overload without spec override. */ public UpdateControl initiateDeployment( BlueGreenContext context, BlueGreenDeploymentType nextBlueGreenDeploymentType, FlinkBlueGreenDeploymentState nextState, Savepoint lastCheckpoint, boolean isFirstDeployment) { + return initiateDeployment( + context, + nextBlueGreenDeploymentType, + nextState, + lastCheckpoint, + isFirstDeployment, + null); + } + + /** + * Initiates a B/G deployment. If {@code specOverride} is non-null, the new child is created + * from it instead of the parent's current spec (used for autoscaler-triggered transitions). + */ + private UpdateControl initiateDeployment( + BlueGreenContext context, + BlueGreenDeploymentType nextBlueGreenDeploymentType, + FlinkBlueGreenDeploymentState nextState, + Savepoint lastCheckpoint, + boolean isFirstDeployment, + @Nullable FlinkBlueGreenDeploymentSpec specOverride) { ObjectMeta bgMeta = context.getBgDeployment().getMetadata(); FlinkDeployment flinkDeployment = @@ -93,7 +110,8 @@ public UpdateControl initiateDeployment( nextBlueGreenDeploymentType, lastCheckpoint, isFirstDeployment, - bgMeta); + bgMeta, + specOverride); deployCluster(context, flinkDeployment); @@ -113,6 +131,7 @@ public UpdateControl initiateDeployment( public UpdateControl checkAndInitiateDeployment( BlueGreenContext context, BlueGreenDeploymentType currentBlueGreenDeploymentType) { + // Check user spec changes first (takes precedence over autoscaler changes) BlueGreenDiffType specDiff = getSpecDiff(context); if (specDiff != BlueGreenDiffType.IGNORE) { @@ -150,7 +169,7 @@ public UpdateControl checkAndInitiateDeployment( savepointTriggered = handleSavepoint(context, currentFlinkDeployment); } catch (Exception e) { var error = "Could not trigger Savepoint. Details: " + e.getMessage(); - return markDeploymentFailing(context, error); + return markDeploymentFailing(context, error, e); } if (savepointTriggered) { @@ -169,7 +188,7 @@ public UpdateControl checkAndInitiateDeployment( } catch (Exception e) { var error = "Could not start Transition. Details: " + e.getMessage(); context.getDeploymentStatus().setSavepointTriggerId(null); - return markDeploymentFailing(context, error); + return markDeploymentFailing(context, error, e); } } else if (specDiff == BlueGreenDiffType.SAVEPOINT_REDEPLOY) { @@ -188,7 +207,7 @@ public UpdateControl checkAndInitiateDeployment( var error = "Could not start Savepoint Redeploy Transition. Details: " + e.getMessage(); - return markDeploymentFailing(context, error); + return markDeploymentFailing(context, error, e); } } else { setLastReconciledSpec(context); @@ -216,6 +235,20 @@ public UpdateControl checkAndInitiateDeployment( } } + // No user spec changes — check if autoscaler produced overrides on the active child. + // If so, trigger a B/G transition with those overrides applied to the new child. + // Skip when FAILING to prevent an infinite loop: abort sets FAILING and rolls back + // to ACTIVE, but the old child still has the same autoscaler overrides → re-detected + // → new transition → abort → re-detected → ... A user spec change clears FAILING + // via the normal transition path above (sets RECONCILING in initiateDeployment). + if (context.getDeploymentStatus().getJobStatus().getState() != JobStatus.FAILING) { + var autoscalerResult = + checkAutoscalerTriggeredTransition(context, currentBlueGreenDeploymentType); + if (autoscalerResult != null) { + return autoscalerResult; + } + } + return UpdateControl.noUpdate(); } @@ -224,9 +257,7 @@ private boolean isChildSuspended(FlinkDeployment deployment) { return false; } var job = deployment.getSpec().getJob(); - return job != null - && job.getState() - == org.apache.flink.kubernetes.operator.api.spec.JobState.SUSPENDED; + return job != null && job.getState() == JobState.SUSPENDED; } private UpdateControl patchFlinkDeployment( @@ -305,8 +336,30 @@ private UpdateControl startTransition( BlueGreenContext context, BlueGreenDeploymentType currentBlueGreenDeploymentType, FlinkDeployment currentFlinkDeployment) { + return startTransition( + context, currentBlueGreenDeploymentType, currentFlinkDeployment, null); + } + + /** + * Starts a B/G transition, optionally using a pre-built spec override for the new child. + * + * @param specOverride if non-null, passed through to {@code prepareFlinkDeployment} so the new + * child is created from this spec instead of the parent's current spec. + */ + private UpdateControl startTransition( + BlueGreenContext context, + BlueGreenDeploymentType currentBlueGreenDeploymentType, + FlinkDeployment currentFlinkDeployment, + @Nullable FlinkBlueGreenDeploymentSpec specOverride) { DeploymentTransition transition = calculateTransition(currentBlueGreenDeploymentType); + LOG.info( + "Starting Blue/Green transition for '{}': {} -> {} (from child '{}')", + context.getDeploymentName(), + currentBlueGreenDeploymentType, + transition.nextBlueGreenDeploymentType, + currentFlinkDeployment.getMetadata().getName()); + Savepoint lastCheckpoint = configureInitialSavepoint(context, currentFlinkDeployment); return initiateDeployment( @@ -314,7 +367,8 @@ private UpdateControl startTransition( transition.nextBlueGreenDeploymentType, transition.nextState, lastCheckpoint, - false); + false, + specOverride); } /** @@ -436,7 +490,7 @@ private boolean handleSavepoint( if (savepointTriggerId == null || savepointTriggerId.isEmpty()) { String triggerId = triggerSavepoint(ctx); - LOG.info("Savepoint requested (triggerId: {}", triggerId); + LOG.info("Savepoint requested (triggerId: {})", triggerId); context.getDeploymentStatus().setSavepointTriggerId(triggerId); return true; } @@ -516,35 +570,32 @@ private UpdateControl finalizeSuspendedDeployment( private UpdateControl handleSpecChangesDuringTransition( BlueGreenContext context, BlueGreenDeploymentType currentBlueGreenDeploymentType) { - if (hasSpecChanged(context)) { - BlueGreenDiffType diffType = getSpecDiff(context); - // Block SUSPEND during transition - wait for transition to complete first - if (diffType == BlueGreenDiffType.SUSPEND) { - LOG.info( - "Suspend requested during transition for '{}'. " - + "Waiting for transition to complete before processing suspend.", - context.getBgDeployment().getMetadata().getName()); - return null; - } + // Parent spec is never persisted with autoscaler values, + // so getSpecDiff detects only user-initiated changes. + BlueGreenDiffType diffType = getSpecDiff(context); - if (diffType != BlueGreenDiffType.IGNORE) { - setLastReconciledSpec(context); - var oppositeDeploymentType = - context.getOppositeDeploymentType(currentBlueGreenDeploymentType); - LOG.info( - "Patching FlinkDeployment '{}' during handleSpecChangesDuringTransition", - context.getDeploymentByType(oppositeDeploymentType) - .getMetadata() - .getName()); - return patchFlinkDeployment( - context, - oppositeDeploymentType, - diffType != BlueGreenDiffType.SAVEPOINT_REDEPLOY); - } + if (diffType == BlueGreenDiffType.IGNORE) { + return null; } - return null; + // Block SUSPEND during transition - wait for transition to complete first + if (diffType == BlueGreenDiffType.SUSPEND) { + LOG.info( + "Suspend requested during transition for '{}'. " + + "Waiting for transition to complete before processing suspend.", + context.getBgDeployment().getMetadata().getName()); + return null; + } + + setLastReconciledSpec(context); + var oppositeDeploymentType = + context.getOppositeDeploymentType(currentBlueGreenDeploymentType); + LOG.info( + "Patching FlinkDeployment '{}' during handleSpecChangesDuringTransition", + context.getDeploymentByType(oppositeDeploymentType).getMetadata().getName()); + return patchFlinkDeployment( + context, oppositeDeploymentType, diffType != BlueGreenDiffType.SAVEPOINT_REDEPLOY); } private TransitionState determineTransitionState( @@ -632,9 +683,11 @@ private UpdateControl deleteDeployment( boolean deleted = deleteFlinkDeployment(currentDeployment, context); if (!deleted) { - LOG.info("FlinkDeployment '{}' not deleted, will retry", currentDeployment); + LOG.warn( + "FlinkDeployment '{}' not deleted, will retry", + currentDeployment.getMetadata().getName()); } else { - LOG.info("FlinkDeployment '{}' deleted!", currentDeployment); + LOG.info("FlinkDeployment '{}' deleted!", currentDeployment.getMetadata().getName()); } return UpdateControl.noUpdate().rescheduleAfter(RETRY_DELAY_MS); @@ -682,6 +735,15 @@ private UpdateControl abortDeployment( suspendFlinkDeployment(context, nextDeployment); + // Clear stale savepointTriggerId to prevent reuse after abort (defense-in-depth). + context.getDeploymentStatus().setSavepointTriggerId(null); + + // Autoscaler overrides persist in the active child's lastReconciledSpec. + // They will NOT be re-detected immediately because markDeploymentFailing sets + // FAILING, and the FAILING guard in checkAndInitiateDeployment blocks autoscaler + // checks — preventing an infinite abort→re-trigger→abort loop. + // A user spec change is required to clear FAILING and unblock transitions. + FlinkBlueGreenDeploymentState previousState = getPreviousState(nextState, context.getDeployments()); context.getDeploymentStatus().setBlueGreenState(previousState); @@ -700,6 +762,13 @@ private static UpdateControl markDeploymentFailing( return patchStatusUpdateControl(context, null, JobStatus.FAILING, error); } + @NotNull + private static UpdateControl markDeploymentFailing( + BlueGreenContext context, String error, Throwable cause) { + LOG.error(error, cause); + return patchStatusUpdateControl(context, null, JobStatus.FAILING, error); + } + private static FlinkBlueGreenDeploymentState getPreviousState( FlinkBlueGreenDeploymentState nextState, FlinkBlueGreenDeployments deployments) { FlinkBlueGreenDeploymentState previousState; @@ -812,6 +881,112 @@ public void updateBlueGreenIngress( blueGreenContext.getJosdkContext()); } + // ==================== Autoscaler Detection for Blue/Green ==================== + + /** + * Detects autoscaler overrides on the active child and initiates a B/G transition if any are + * found. A merged spec copy (parent spec + overrides) is built and passed explicitly to the + * transition flow so the new child inherits the overrides without mutating the parent's + * in-memory spec. + * + * @return UpdateControl if a transition was initiated, or {@code null} if no overrides + */ + @Nullable + private UpdateControl checkAutoscalerTriggeredTransition( + BlueGreenContext context, BlueGreenDeploymentType currentType) { + + FlinkDeployment activeChild = context.getDeploymentByType(currentType); + if (activeChild == null || !isChildStableForAutoscalerPropagation(activeChild)) { + return null; + } + + // Stability guarantees reconStatus is non-null and not before first deployment. + FlinkDeploymentSpec childLastReconciledSpec = + activeChild.getStatus().getReconciliationStatus().deserializeLastReconciledSpec(); + if (childLastReconciledSpec == null) { + return null; + } + + // Reuse the same ReflectiveDiffBuilder as the child reconciler — any field the + // autoscaler touches is detected automatically, including future capabilities. + var autoscalerDiff = + new ReflectiveDiffBuilder<>( + KubernetesDeploymentMode.getDeploymentMode(activeChild), + activeChild.getSpec(), + childLastReconciledSpec) + .build(); + + if (autoscalerDiff.getType() == DiffType.IGNORE) { + return null; + } + + LOG.info( + "Autoscaler overrides detected on child '{}' (diffType: {}). " + + "Initiating Blue/Green transition.", + activeChild.getMetadata().getName(), + autoscalerDiff.getType()); + + // Savepoint first if required. Don't setLastReconciledSpec yet — overrides + // must remain detectable from lastReconciledSpec after savepoint completes. + try { + boolean savepointTriggered = handleSavepoint(context, activeChild); + if (savepointTriggered) { + var savepointingState = calculateSavepointingState(currentType); + return patchStatusUpdateControl(context, savepointingState, null, null) + .rescheduleAfter(getReconciliationReschedInterval(context)); + } + } catch (Exception e) { + var error = + "Could not trigger Savepoint for autoscaler transition. Details: " + + e.getMessage(); + return markDeploymentFailing(context, error, e); + } + + setLastReconciledSpec(context); + FlinkBlueGreenDeploymentSpec mergedSpec = buildMergedSpec(context, childLastReconciledSpec); + + try { + return startTransition(context, currentType, activeChild, mergedSpec); + } catch (Exception e) { + var error = "Could not start autoscaler transition. Details: " + e.getMessage(); + context.getDeploymentStatus().setSavepointTriggerId(null); + return markDeploymentFailing(context, error, e); + } + } + + /** + * Deep-copies the parent spec and overwrites the autoscaler-managed fields ({@code + * flinkConfiguration}, {@code taskManager}) wholesale from the child's {@code + * lastReconciledSpec}. We don't replace the entire {@code template.spec} because it contains + * deployment-specific transforms (e.g., ingress prefixing). + */ + private FlinkBlueGreenDeploymentSpec buildMergedSpec( + BlueGreenContext context, FlinkDeploymentSpec childLastReconciledSpec) { + FlinkBlueGreenDeploymentSpec merged = + SpecUtils.readSpecFromJSON( + SpecUtils.writeSpecAsJSON(context.getBgDeployment().getSpec(), "spec"), + "spec", + FlinkBlueGreenDeploymentSpec.class); + var parentFlinkSpec = merged.getTemplate().getSpec(); + parentFlinkSpec.setFlinkConfiguration(childLastReconciledSpec.getFlinkConfiguration()); + parentFlinkSpec.setTaskManager(childLastReconciledSpec.getTaskManager()); + return merged; + } + + /** + * Checks if a child FlinkDeployment has been stable long enough to trust autoscaler changes. + * Requires STABLE lifecycle state and lastReconciledSpecStable to prevent cascading + * transitions. + */ + private boolean isChildStableForAutoscalerPropagation(FlinkDeployment child) { + if (child.getStatus() == null + || child.getStatus().getLifecycleState() != ResourceLifecycleState.STABLE) { + return false; + } + var reconStatus = child.getStatus().getReconciliationStatus(); + return reconStatus != null && reconStatus.isLastReconciledSpecStable(); + } + // ==================== Common Utility Methods ==================== public static UpdateControl patchStatusUpdateControl( @@ -839,7 +1014,7 @@ public static UpdateControl patchStatusUpdateControl( deploymentStatus.setError(null); } - deploymentStatus.setLastReconciledTimestamp(java.time.Instant.now().toString()); + deploymentStatus.setLastReconciledTimestamp(Instant.now().toString()); flinkBlueGreenDeployment.setStatus(deploymentStatus); return UpdateControl.patchStatus(flinkBlueGreenDeployment); } diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenKubernetesService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenKubernetesService.java index ae7492d312..2f64b1ab35 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenKubernetesService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/BlueGreenKubernetesService.java @@ -116,4 +116,29 @@ public static boolean deleteFlinkDeployment( return deletedStatus.size() == 1 && deletedStatus.get(0).getKind().equals("FlinkDeployment"); } + + /** + * Checks if a FlinkDeployment is owned by a FlinkBlueGreenDeployment. + * + *

This is used to determine whether in-place scaling should be skipped for this deployment, + * as Blue/Green deployments handle scaling through full transitions instead. + * + * @param deployment the FlinkDeployment to check + * @return true if the deployment is owned by a FlinkBlueGreenDeployment, false otherwise + */ + public static boolean isOwnedByBlueGreenDeployment(FlinkDeployment deployment) { + if (deployment == null || deployment.getMetadata() == null) { + return false; + } + var ownerReferences = deployment.getMetadata().getOwnerReferences(); + if (ownerReferences == null || ownerReferences.isEmpty()) { + return false; + } + return ownerReferences.stream() + .anyMatch( + ref -> + FlinkBlueGreenDeployment.class + .getSimpleName() + .equals(ref.getKind())); + } } diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractFlinkResourceReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractFlinkResourceReconciler.java index 7f8af3ca88..ae1b56f6a0 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractFlinkResourceReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractFlinkResourceReconciler.java @@ -34,6 +34,7 @@ import org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext; import org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions; import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext; +import org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenKubernetesService; import org.apache.flink.kubernetes.operator.reconciler.Reconciler; import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils; import org.apache.flink.kubernetes.operator.reconciler.diff.DiffResult; @@ -129,6 +130,22 @@ public void reconcile(FlinkResourceContext ctx) throws Exception { cr.getStatus().getReconciliationStatus().deserializeLastReconciledSpec(); SPEC currentDeploySpec = cr.getSpec(); + // Snapshot pre-autoscaler diff for B/G children so we can distinguish + // autoscaler-only changes from parent-initiated ones after applyAutoscaler(). + boolean isBlueGreenOwned = + cr instanceof FlinkDeployment + && BlueGreenKubernetesService.isOwnedByBlueGreenDeployment( + (FlinkDeployment) cr); + boolean hasPreAutoscalerSpecChange = + isBlueGreenOwned + && DiffType.IGNORE + != new ReflectiveDiffBuilder<>( + ctx.getDeploymentMode(), + lastReconciledSpec, + currentDeploySpec) + .build() + .getType(); + applyAutoscaler(ctx); var reconciliationState = reconciliationStatus.getState(); @@ -152,6 +169,19 @@ public void reconcile(FlinkResourceContext ctx) throws Exception { if (checkNewSpecAlreadyDeployed(cr, deployConfig)) { return; } + + // Autoscaler-only diff on a B/G child → skip in-place processing. + // The B/G parent will detect this via lastReconciledSpec and trigger a transition. + if (isBlueGreenOwned && !hasPreAutoscalerSpecChange) { + LOG.info( + "Deferring autoscaler change ({}) on B/G child '{}' to parent transition.", + diffType, + cr.getMetadata().getName()); + ReconciliationUtils.updateStatusForDeployedSpec( + ctx.getResource(), deployConfig, clock); + return; + } + triggerSpecChangeEvent(cr, specDiff, ctx.getKubernetesClient()); // Try scaling if this is not an upgrade/redeploy change 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..dd3d837abc 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 @@ -27,6 +27,7 @@ import org.apache.flink.kubernetes.operator.api.bluegreen.BlueGreenDiffType; import org.apache.flink.kubernetes.operator.api.spec.FlinkBlueGreenDeploymentSpec; import org.apache.flink.kubernetes.operator.api.spec.IngressSpec; +import org.apache.flink.kubernetes.operator.api.spec.JobState; import org.apache.flink.kubernetes.operator.api.spec.KubernetesDeploymentMode; import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode; import org.apache.flink.kubernetes.operator.api.status.FlinkBlueGreenDeploymentStatus; @@ -40,6 +41,7 @@ import org.apache.flink.util.Preconditions; import io.fabric8.kubernetes.api.model.ObjectMeta; +import org.jetbrains.annotations.Nullable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -179,9 +181,7 @@ public static long getDeploymentDeletionDelay(BlueGreenContext context) { * @return abort grace period in milliseconds */ public static long getAbortGracePeriod(BlueGreenContext context) { - long abortGracePeriod = - getConfigOption(context.getBgDeployment(), ABORT_GRACE_PERIOD).toMillis(); - return abortGracePeriod; + return getConfigOption(context.getBgDeployment(), ABORT_GRACE_PERIOD).toMillis(); } /** @@ -213,9 +213,7 @@ public static boolean isSavepointRequired(BlueGreenContext context) { .getSpec() .getJob() .getUpgradeMode(); - // return UpgradeMode.SAVEPOINT == upgradeMode; // Currently taking savepoints for all modes except STATELESS - // (previously only SAVEPOINT mode required savepoints) return UpgradeMode.STATELESS != upgradeMode; } @@ -302,26 +300,33 @@ public static Savepoint getLastCheckpoint( // ==================== Deployment Preparation Utilities ==================== + /** Convenience overload without spec override. */ + public static FlinkDeployment prepareFlinkDeployment( + BlueGreenContext context, + BlueGreenDeploymentType blueGreenDeploymentType, + Savepoint lastCheckpoint, + boolean isFirstDeployment, + ObjectMeta bgMeta) { + return prepareFlinkDeployment( + context, blueGreenDeploymentType, lastCheckpoint, isFirstDeployment, bgMeta, null); + } + /** - * Creates a new FlinkDeployment resource for a Blue/Green deployment transition. This method - * prepares the deployment with proper metadata, specs, and savepoint configuration. - * - * @param context the Blue/Green transition context - * @param blueGreenDeploymentType the type of deployment (BLUE or GREEN) - * @param lastCheckpoint the savepoint/checkpoint to restore from (can be null) - * @param isFirstDeployment whether this is the initial deployment - * @param bgMeta the metadata of the parent Blue/Green deployment - * @return configured FlinkDeployment ready for deployment + * Creates a FlinkDeployment for a B/G transition. If {@code specOverride} is non-null, it is + * used as the source spec instead of the parent's current spec (e.g., for autoscaler + * overrides). */ public static FlinkDeployment prepareFlinkDeployment( BlueGreenContext context, BlueGreenDeploymentType blueGreenDeploymentType, Savepoint lastCheckpoint, boolean isFirstDeployment, - ObjectMeta bgMeta) { + ObjectMeta bgMeta, + @Nullable FlinkBlueGreenDeploymentSpec specOverride) { // Deployment FlinkDeployment flinkDeployment = new FlinkDeployment(); - FlinkBlueGreenDeploymentSpec originalSpec = context.getBgDeployment().getSpec(); + FlinkBlueGreenDeploymentSpec originalSpec = + specOverride != null ? specOverride : context.getBgDeployment().getSpec(); String childDeploymentName = bgMeta.getName() + "-" + blueGreenDeploymentType.toString().toLowerCase(); @@ -343,14 +348,13 @@ public static FlinkDeployment prepareFlinkDeployment( String initialSavepointPath = spec.getTemplate().getSpec().getJob().getInitialSavepointPath(); if (initialSavepointPath != null && !initialSavepointPath.isEmpty()) { - LOG.info("Using initialSavepointPath: " + initialSavepointPath); - spec.getTemplate().getSpec().getJob().setInitialSavepointPath(initialSavepointPath); + LOG.info("Using initialSavepointPath: {}", initialSavepointPath); } else { LOG.info("Clean startup with no checkpoint/savepoint restoration"); } } else if (lastCheckpoint != null) { String location = lastCheckpoint.getLocation().replace("file:", ""); - LOG.info("Using Blue/Green savepoint/checkpoint: " + location); + LOG.info("Using Blue/Green savepoint/checkpoint: {}", location); spec.getTemplate().getSpec().getJob().setInitialSavepointPath(location); } else { String initialSavepointPath = @@ -366,6 +370,16 @@ public static FlinkDeployment prepareFlinkDeployment( flinkDeployment.setSpec(spec.getTemplate().getSpec()); + // Ensure job.state is never null (which could cause merge issues downstream). + // The spec was already copied on the line above, so this only needs to handle the + // null-state defensive case. + if (flinkDeployment.getSpec().getJob().getState() == null) { + flinkDeployment.getSpec().getJob().setState(JobState.RUNNING); + LOG.warn( + "Job state was null on parent spec for '{}', defaulting to RUNNING", + flinkDeployment.getMetadata().getName()); + } + // Update Ingress template if exists to prevent path collision between Blue and Green IngressSpec ingress = flinkDeployment.getSpec().getIngress(); if (ingress != null) { 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 528b61727d..c462dfdd51 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 @@ -20,6 +20,7 @@ import org.apache.flink.api.common.JobStatus; import org.apache.flink.configuration.CheckpointingOptions; import org.apache.flink.configuration.Configuration; +import org.apache.flink.configuration.PipelineOptions; import org.apache.flink.configuration.TaskManagerOptions; import org.apache.flink.kubernetes.operator.TestUtils; import org.apache.flink.kubernetes.operator.TestingFlinkService; @@ -62,6 +63,7 @@ import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.function.Consumer; import java.util.stream.Stream; import static org.apache.flink.kubernetes.operator.api.spec.FlinkBlueGreenDeploymentConfigOptions.ABORT_GRACE_PERIOD; @@ -257,6 +259,372 @@ public void verifySavepointRedeployNonceTriggersTransitionWithInitialSavepointPa rs.reconciledStatus.getBlueGreenState()); } + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyAutoscalerChangesTriggersBlueGreenTransition(FlinkVersion flinkVersion) + throws Exception { + var rs = + executeBasicDeployment( + flinkVersion, buildAutoscalerEnabledCluster(flinkVersion), false, null); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, rs.reconciledStatus.getBlueGreenState()); + + // Simulate autoscaler adding parallelism overrides to the active Blue child + String autoscalerOverrides = "vertex1:4,vertex2:8"; + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> + spec.getFlinkConfiguration() + .put( + PipelineOptions.PARALLELISM_OVERRIDES.key(), + autoscalerOverrides)); + + // Reconcile — should trigger a B/G transition + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "Autoscaler changes should trigger transition to Green"); + + // Verify: new Green child has the overrides (parent spec is NOT modified) + assertEquals( + autoscalerOverrides, + getChildByColor("green") + .getSpec() + .getFlinkConfiguration() + .asFlatMap() + .get(PipelineOptions.PARALLELISM_OVERRIDES.key()), + "Green deployment should have the parallelism overrides"); + + // Complete the transition + rs = completeTransition(rs, "green"); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "Should be ACTIVE_GREEN after transition completes"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyVerticalAutoscalerChangesTriggerBlueGreenTransition(FlinkVersion flinkVersion) + throws Exception { + // Setup: Deploy with initial TM memory and autoscaler enabled + var deployment = buildAutoscalerEnabledCluster(flinkVersion); + var tmSpec = new TaskManagerSpec(); + var tmResource = new Resource(); + tmResource.setMemory("2048m"); + tmSpec.setResource(tmResource); + deployment.getSpec().getTemplate().getSpec().setTaskManager(tmSpec); + + var rs = executeBasicDeployment(flinkVersion, deployment, false, null); + + // Simulate autoscaler performing vertical scaling (memory tuning) + String scaledMemory = String.valueOf(4096 * 1024 * 1024L); + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> { + spec.getFlinkConfiguration() + .put(TaskManagerOptions.TOTAL_PROCESS_MEMORY.key(), scaledMemory); + spec.getFlinkConfiguration() + .put(TaskManagerOptions.FRAMEWORK_HEAP_MEMORY.key(), "0 bytes"); + spec.getTaskManager().getResource().setMemory(scaledMemory); + }); + + // Reconcile — should trigger a B/G transition + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "Vertical autoscaler changes should trigger transition to Green"); + + // Verify Green child has the updated memory + var greenChild = getChildByColor("green"); + assertEquals( + scaledMemory, + greenChild.getSpec().getTaskManager().getResource().getMemory(), + "Green deployment should have the autoscaler's TM memory"); + assertNotNull( + greenChild + .getSpec() + .getFlinkConfiguration() + .get(TaskManagerOptions.TOTAL_PROCESS_MEMORY.key()), + "Green deployment should have TM total process memory config"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyUserParallelismChangeAfterAutoscalerTriggersTransition( + FlinkVersion flinkVersion) throws Exception { + var rs = + executeBasicDeployment( + flinkVersion, buildAutoscalerEnabledCluster(flinkVersion), false, null); + + // Autoscaler adds parallelism overrides → B/G transition → complete to ACTIVE_GREEN + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> + spec.getFlinkConfiguration() + .put( + PipelineOptions.PARALLELISM_OVERRIDES.key(), + "vertex1:4,vertex2:8")); + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState()); + rs = completeTransition(rs, "green"); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState()); + + // User manually sets parallelism overrides → should trigger a new B/G transition. + rs.deployment = kubernetesClient.resource(rs.deployment).get(); + rs.deployment + .getSpec() + .getTemplate() + .getSpec() + .getFlinkConfiguration() + .put(PipelineOptions.PARALLELISM_OVERRIDES.key(), "vertex1:2,vertex2:4"); + kubernetesClient.resource(rs.deployment).createOrReplace(); + + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "User parallelism change after autoscaler should trigger a new transition"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyUserTmMemoryChangeAfterAutoscalerTriggersTransition(FlinkVersion flinkVersion) + throws Exception { + var deployment = buildAutoscalerEnabledCluster(flinkVersion); + deployment + .getSpec() + .getTemplate() + .getSpec() + .getFlinkConfiguration() + .put("job.autoscaler.memory.tuning.enabled", "true"); + var tmSpec = new TaskManagerSpec(); + var tmResource = new Resource(); + tmResource.setMemory("2048m"); + tmSpec.setResource(tmResource); + deployment.getSpec().getTemplate().getSpec().setTaskManager(tmSpec); + + var rs = executeBasicDeployment(flinkVersion, deployment, false, null); + + // Autoscaler changes TM memory → B/G transition → complete to ACTIVE_GREEN + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> spec.getTaskManager().getResource().setMemory("4096m")); + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState()); + rs = completeTransition(rs, "green"); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState()); + + // User manually changes TM memory → should trigger a new B/G transition. + rs.deployment = kubernetesClient.resource(rs.deployment).get(); + rs.deployment + .getSpec() + .getTemplate() + .getSpec() + .getTaskManager() + .getResource() + .setMemory("3072m"); + kubernetesClient.resource(rs.deployment).createOrReplace(); + + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "User TM memory change after autoscaler should trigger a new transition"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyAutoscalerPhantomDiffPreventionAfterTransition(FlinkVersion flinkVersion) + throws Exception { + var rs = + executeBasicDeployment( + flinkVersion, buildAutoscalerEnabledCluster(flinkVersion), false, null); + + // Autoscaler adds parallelism overrides → B/G transition → complete + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> + spec.getFlinkConfiguration() + .put( + PipelineOptions.PARALLELISM_OVERRIDES.key(), + "vertex1:4,vertex2:8")); + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState()); + rs = completeTransition(rs, "green"); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState()); + + // Subsequent reconciles should NOT trigger spurious transitions. + // Parent spec and lastReconciledSpec were never modified with autoscaler values, + // so there is no phantom diff. The green child was created with overrides baked into + // its K8s spec, and its lastReconciledSpec matches (no further autoscaler diff). + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "No phantom diff — should remain ACTIVE_GREEN"); + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "Should still be ACTIVE_GREEN after multiple reconciles"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyAutoscalerDoesNotRetriggerAfterAbortedTransition(FlinkVersion flinkVersion) + throws Exception { + var rs = + executeBasicDeployment( + flinkVersion, buildAutoscalerEnabledCluster(flinkVersion), false, null); + + // Autoscaler adds parallelism overrides → triggers B/G transition + String autoscalerOverrides = "vertex1:4,vertex2:8"; + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> + spec.getFlinkConfiguration() + .put( + PipelineOptions.PARALLELISM_OVERRIDES.key(), + autoscalerOverrides)); + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState()); + + // Green child fails to start — abort after grace period. + Thread.sleep(MINIMUM_ABORT_GRACE_PERIOD); + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should roll back to ACTIVE_BLUE after abort"); + + // Clean up stale Green child. + try { + kubernetesClient.resource(getChildByColor("green")).delete(); + } catch (AssertionError ignored) { + // Green may already be deleted + } + + // After abort the parent is in FAILING state. The FAILING guard prevents the + // autoscaler from re-triggering the same overrides that just failed — this avoids + // an infinite abort→re-trigger→abort loop. A user spec change is required to + // clear FAILING and unblock new transitions. + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should remain ACTIVE_BLUE (FAILING) — autoscaler must not re-trigger after abort"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyAutoscalerTransitionWithSavepointUpgradeMode(FlinkVersion flinkVersion) + throws Exception { + // Setup: SAVEPOINT upgrade mode with autoscaler enabled + var deployment = + buildSessionCluster( + TEST_DEPLOYMENT_NAME, + TEST_NAMESPACE, + flinkVersion, + null, + UpgradeMode.SAVEPOINT); + deployment + .getSpec() + .getTemplate() + .getSpec() + .getFlinkConfiguration() + .put("job.autoscaler.enabled", "true"); + + var rs = executeBasicDeployment(flinkVersion, deployment, false, null); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, rs.reconciledStatus.getBlueGreenState()); + + // Simulate autoscaler adding parallelism overrides to the active Blue child + String autoscalerOverrides = "vertex1:4,vertex2:8"; + simulateAutoscalerOnChild( + getChildByColor("blue"), + spec -> + spec.getFlinkConfiguration() + .put( + PipelineOptions.PARALLELISM_OVERRIDES.key(), + autoscalerOverrides)); + + // Drive the savepoint flow — reuse the existing handleSavepoint helper. + // First reconcile detects overrides and triggers a savepoint → SAVEPOINTING_BLUE + var triggers = flinkService.getSavepointTriggers(); + triggers.clear(); + + rs = reconcile(rs.deployment); + + // Simulating a pending savepoint + triggers.put(rs.deployment.getStatus().getSavepointTriggerId(), false); + + assertEquals( + FlinkBlueGreenDeploymentState.SAVEPOINTING_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Autoscaler transition should enter SAVEPOINTING_BLUE for SAVEPOINT upgrade mode"); + assertTrue(rs.updateControl.isPatchStatus()); + + // Next reconciliation waits on the pending savepoint + rs = reconcile(rs.deployment); + assertTrue(rs.updateControl.isNoUpdate()); + + // Complete the savepoint + triggers.put(rs.deployment.getStatus().getSavepointTriggerId(), true); + rs = reconcile(rs.deployment); + + // Should return to ACTIVE_BLUE after savepoint completes + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should return to ACTIVE_BLUE after savepoint completes"); + + // Next reconcile re-detects autoscaler overrides and starts the actual transition + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "Should transition to Green after savepoint completes"); + + // Verify: Green child has the autoscaler overrides + assertEquals( + autoscalerOverrides, + getChildByColor("green") + .getSpec() + .getFlinkConfiguration() + .asFlatMap() + .get(PipelineOptions.PARALLELISM_OVERRIDES.key()), + "Green deployment should have the parallelism overrides"); + + // Verify: Green child uses the savepoint (not null) + assertNotNull( + getChildByColor("green").getSpec().getJob().getInitialSavepointPath(), + "Green deployment should use a savepoint path"); + + // Complete the transition + rs = completeTransition(rs, "green"); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState(), + "Should be ACTIVE_GREEN after transition completes"); + } + @ParameterizedTest @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") public void verifySuspendAndResumeInPlace(FlinkVersion flinkVersion) throws Exception { @@ -1107,6 +1475,64 @@ static class ReconcileResult { flinkVersion, blueGreenDeployment, false, TEST_INITIAL_SAVEPOINT_PATH); } + // ---- Autoscaler Test Helpers ---- + + /** Builds a session cluster with the autoscaler enabled in the parent's flinkConfiguration. */ + private FlinkBlueGreenDeployment buildAutoscalerEnabledCluster(FlinkVersion flinkVersion) { + var deployment = + buildSessionCluster( + TEST_DEPLOYMENT_NAME, + TEST_NAMESPACE, + flinkVersion, + null, + UpgradeMode.STATELESS); + deployment + .getSpec() + .getTemplate() + .getSpec() + .getFlinkConfiguration() + .put("job.autoscaler.enabled", "true"); + return deployment; + } + + /** + * Simulates the autoscaler modifying the child's in-memory spec. Applies the given mutation to + * the child's spec, serializes it into lastReconciledSpec, marks it stable, and updates status. + */ + private void simulateAutoscalerOnChild( + FlinkDeployment child, Consumer specMutation) { + var spec = child.getSpec(); + specMutation.accept(spec); + child.getStatus().getReconciliationStatus().serializeAndSetLastReconciledSpec(spec, child); + child.getStatus().getReconciliationStatus().markReconciledSpecAsStable(); + kubernetesClient.resource(child).updateStatus(); + } + + /** Finds a child FlinkDeployment by color suffix ("blue" or "green"). */ + private FlinkDeployment getChildByColor(String color) { + return getFlinkDeployments().stream() + .filter(d -> d.getMetadata().getName().endsWith("-" + color)) + .findFirst() + .orElseThrow(() -> new AssertionError("No " + color + " child found")); + } + + /** + * Drives a B/G transition to completion: simulates a successful job start on the target child, + * waits for the deletion delay, and reconciles until the new color is active. + */ + private TestingFlinkBlueGreenDeploymentController.BlueGreenReconciliationResult + completeTransition( + TestingFlinkBlueGreenDeploymentController.BlueGreenReconciliationResult rs, + String targetColor) + throws Exception { + simulateSuccessfulJobStart(getChildByColor(targetColor)); + rs = reconcile(rs.deployment); + Thread.sleep(rs.updateControl.getScheduleDelay().get()); + reconcile(rs.deployment); + assertEquals(1, getFlinkDeployments().size()); + return reconcile(rs.deployment); + } + private ReconcileResult reconcileAndVerifyPatchBehavior( TestingFlinkBlueGreenDeploymentController.BlueGreenReconciliationResult rs) throws Exception {