From 5db841caef457732529d6685867b47459e765503 Mon Sep 17 00:00:00 2001 From: Ihor Mielientiev Date: Thu, 26 Feb 2026 21:45:24 +0100 Subject: [PATCH 1/3] [FLINK-39165] Prevent concurrent FlinkSessionJob / FlinkDeployment update --- .../deployment/SessionReconciler.java | 35 +++++ .../sessionjob/SessionJobReconciler.java | 30 +++- .../flink/kubernetes/operator/TestUtils.java | 22 ++- .../deployment/SessionReconcilerTest.java | 69 +++++++++ .../sessionjob/SessionJobReconcilerTest.java | 134 ++++++++++++++++++ 5 files changed, 286 insertions(+), 4 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java index 4c998fec16..bc86f37733 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java @@ -27,6 +27,7 @@ import org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec; import org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus; import org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus; +import org.apache.flink.kubernetes.operator.api.status.ReconciliationState; import org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions; import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext; import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils; @@ -42,6 +43,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.Objects; import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; @@ -64,9 +66,42 @@ public SessionReconciler( @Override protected boolean readyToReconcile(FlinkResourceContext ctx) { + var deployment = ctx.getResource(); + // Only check when there is a pending spec change on the deployment itself + if (!Objects.equals( + deployment.getMetadata().getGeneration(), + deployment.getStatus().getObservedGeneration())) { + boolean hasUpgradingSessionJob = + ctx.getJosdkContext().getSecondaryResources(FlinkSessionJob.class).stream() + .filter( + job -> + deployment + .getMetadata() + .getName() + .equals(job.getSpec().getDeploymentName())) + .anyMatch(SessionReconciler::isSessionJobTransitioning); + if (hasUpgradingSessionJob) { + LOG.info( + "Deferring session cluster upgrade: associated session jobs are being upgraded"); + return false; + } + } return true; } + private static boolean isSessionJobTransitioning(FlinkSessionJob job) { + var reconStatus = job.getStatus().getReconciliationStatus(); + // New session jobs state is UPGRADING by default before first deployment. Treating them as + // transitioning would deadlock: the deployment can't + // start because it waits for the job, and the job can't start because it waits for the + // cluster. + if (reconStatus.isBeforeFirstDeployment()) { + return false; + } + var state = reconStatus.getState(); + return state == ReconciliationState.UPGRADING || state == ReconciliationState.ROLLING_BACK; + } + @Override protected boolean reconcileSpecChange( DiffType diffType, diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java index f5d10626b1..f70a65841b 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java @@ -27,6 +27,7 @@ import org.apache.flink.kubernetes.operator.api.spec.JobState; import org.apache.flink.kubernetes.operator.api.status.FlinkSessionJobStatus; import org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus; +import org.apache.flink.kubernetes.operator.api.status.ReconciliationState; import org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext; import org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions; import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext; @@ -43,6 +44,7 @@ import org.slf4j.LoggerFactory; import java.time.Instant; +import java.util.Objects; import java.util.Optional; /** The reconciler for the {@link FlinkSessionJob}. */ @@ -186,9 +188,33 @@ public static boolean sessionClusterReady(Optional flinkDeploym "Session cluster deployment is in {} status, not ready for serve", jobmanagerDeploymentStatus); return false; - } else { - return true; } + + // Block while FlinkDeployment is in reconciliation state + var reconciliationState = flinkdep.getStatus().getReconciliationStatus().getState(); + if (reconciliationState != ReconciliationState.DEPLOYED + && reconciliationState != ReconciliationState.ROLLED_BACK) { + LOG.info( + "Session cluster deployment reconciliation state is {}, not ready for sync", + reconciliationState); + return false; + } + + // Block while a FlinkDeployment spec change is pending. This make sure that + // FlinkSessionJob updates will be + // applied only after all changes for this generation got applied + if (!Objects.equals( + flinkdep.getMetadata().getGeneration(), + flinkdep.getStatus().getObservedGeneration())) { + LOG.info( + "Session cluster deployment has pending spec changes " + + "(generation={}, observedGeneration={}), not ready for sync", + flinkdep.getMetadata().getGeneration(), + flinkdep.getStatus().getObservedGeneration()); + return false; + } + + return true; } else { LOG.warn("Session cluster deployment is not found"); return false; diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestUtils.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestUtils.java index a48f41d9c9..49a682b84b 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestUtils.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestUtils.java @@ -26,6 +26,7 @@ import org.apache.flink.kubernetes.operator.api.spec.JobReference; import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode; import org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus; +import org.apache.flink.kubernetes.operator.api.status.ReconciliationState; import org.apache.flink.kubernetes.operator.api.utils.BaseTestUtils; import org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions; import org.apache.flink.kubernetes.operator.health.CanaryResourceManager; @@ -226,6 +227,12 @@ public Optional getSecondaryResource(Class expectedType, String eventSourceNa session.getStatus() .getReconciliationStatus() .serializeAndSetLastReconciledSpec(session.getSpec(), session); + // ReconciliationStatus defaults to UPGRADING. Ready session + // cluster must be in DEPLOYED state, otherwise sessionClusterReady() will + // reject it due to the reconciliation state check. + session.getStatus() + .getReconciliationStatus() + .setState(ReconciliationState.DEPLOYED); return (Optional) Optional.of(session); } @@ -498,8 +505,11 @@ public Optional getRetryInfo() { @Override public Set getSecondaryResources(Class aClass) { // TODO: improve this, even if we only support FlinkDeployment as a secondary resource + KubernetesClient client = getClient(); + if (client == null) { + return Set.of(); + } if (aClass.getSimpleName().equals(FlinkDeployment.class.getSimpleName())) { - KubernetesClient client = getClient(); var hasMetadata = new HashSet<>( client.resources(FlinkDeployment.class) @@ -507,8 +517,16 @@ public Set getSecondaryResources(Class aClass) { .list() .getItems()); return (Set) hasMetadata; + } else if (aClass.getSimpleName().equals(FlinkSessionJob.class.getSimpleName())) { + var hasMetadata = + new HashSet<>( + client.resources(FlinkSessionJob.class) + .inAnyNamespace() + .list() + .getItems()); + return (Set) hasMetadata; } else { - return null; + return Set.of(); } } diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java index 930ad9a4d5..37028bdd6c 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java @@ -305,4 +305,73 @@ public void testGetNonTerminalJobs() throws Exception { nonTerminalJobsAfterRemoval.size(), "Should have no non-terminal jobs when only terminated jobs exist"); } + + @Test + public void testUpgradeDeferredWhenSessionJobUpgrading() throws Exception { + FlinkDeployment deployment = TestUtils.buildSessionCluster(); + reconciler.reconcile(deployment, flinkService.getContext()); + assertEquals( + ReconciliationState.DEPLOYED, + deployment.getStatus().getReconciliationStatus().getState()); + + // Create a session job that is in UPGRADING state + FlinkSessionJob upgradingJob = TestUtils.buildSessionJob(); + upgradingJob + .getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(upgradingJob.getSpec(), upgradingJob); + upgradingJob.getStatus().getReconciliationStatus().setState(ReconciliationState.UPGRADING); + kubernetesClient.resource(upgradingJob).createOrReplace(); + + // Simulate a pending spec change on the deployment + deployment.getMetadata().setGeneration(2L); + deployment.getSpec().setRestartNonce(2L); + + // Reconcile should be deferred because session job is upgrading + reconciler.reconcile(deployment, flinkService.getContext()); + + // Should still be DEPLOYED (not UPGRADING) — upgrade was deferred + assertEquals( + ReconciliationState.DEPLOYED, + deployment.getStatus().getReconciliationStatus().getState()); + } + + @Test + public void testUpgradeProceedsWhenSessionJobDeployed() throws Exception { + FlinkDeployment deployment = TestUtils.buildSessionCluster(); + + // Initial deploy + reconciler.reconcile(deployment, flinkService.getContext()); + assertEquals( + ReconciliationState.DEPLOYED, + deployment.getStatus().getReconciliationStatus().getState()); + + // Create a session job that is in DEPLOYED state (not upgrading) + FlinkSessionJob deployedJob = TestUtils.buildSessionJob(); + deployedJob + .getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(deployedJob.getSpec(), deployedJob); + deployedJob.getStatus().getReconciliationStatus().setState(ReconciliationState.DEPLOYED); + kubernetesClient.resource(deployedJob).createOrReplace(); + + // Simulate a pending spec change on the deployment + deployment.getMetadata().setGeneration(2L); + deployment.getSpec().setRestartNonce(456L); + + // Reconcile should proceed because session job is not upgrading + reconciler.reconcile(deployment, flinkService.getContext()); + + // Should be DEPLOYED with the new spec (upgrade went through) + assertEquals( + ReconciliationState.DEPLOYED, + deployment.getStatus().getReconciliationStatus().getState()); + assertEquals( + 456L, + deployment + .getStatus() + .getReconciliationStatus() + .deserializeLastReconciledSpec() + .getRestartNonce()); + } } diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java index b001ed7739..e0370a4198 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java @@ -29,6 +29,7 @@ import org.apache.flink.kubernetes.operator.api.spec.JobState; import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode; import org.apache.flink.kubernetes.operator.api.status.FlinkSessionJobStatus; +import org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus; import org.apache.flink.kubernetes.operator.api.status.JobStatus; import org.apache.flink.kubernetes.operator.api.status.ReconciliationState; import org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType; @@ -40,8 +41,10 @@ import org.apache.flink.kubernetes.operator.utils.SnapshotUtils; import org.apache.flink.runtime.client.JobStatusMessage; +import io.fabric8.kubernetes.api.model.HasMetadata; import io.fabric8.kubernetes.client.KubernetesClient; import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient; +import io.javaoperatorsdk.operator.api.reconciler.Context; import lombok.Getter; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -54,6 +57,7 @@ import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ThreadLocalRandom; import java.util.stream.Stream; @@ -861,4 +865,134 @@ public void testJobIdGeneration() throws Exception { // New jobID recorded despite failure Assertions.assertNotEquals(jobID, sessionJob.getStatus().getJobStatus().getJobId()); } + + @Test + public void testSessionClusterNotReadyWhenUpgrading() throws Exception { + FlinkSessionJob sessionJob = TestUtils.buildSessionJob(); + + // Create a context where the deployment is JM=READY but reconciliation state=UPGRADING + Context upgradingContext = + new TestUtils.TestingContext<>() { + @Override + public Optional getSecondaryResource( + Class expectedType, String eventSourceName) { + var session = TestUtils.buildSessionCluster(); + session.getStatus() + .setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); + session.getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(session.getSpec(), session); + // Leave state as UPGRADING (the default) + return Optional.of(session); + } + + @Override + public KubernetesClient getClient() { + return kubernetesClient; + } + }; + + reconciler.reconcile(sessionJob, upgradingContext); + assertEquals(0, flinkService.listJobs().size()); + } + + @Test + public void testSessionClusterNotReadyWhenGenerationMismatch() throws Exception { + FlinkSessionJob sessionJob = TestUtils.buildSessionJob(); + + // Create a context where JM=READY, state=DEPLOYED, but generation != observedGeneration + Context mismatchContext = + new TestUtils.TestingContext<>() { + @Override + public Optional getSecondaryResource( + Class expectedType, String eventSourceName) { + var session = TestUtils.buildSessionCluster(); + session.getStatus() + .setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); + session.getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(session.getSpec(), session); + session.getStatus() + .getReconciliationStatus() + .setState(ReconciliationState.DEPLOYED); + // Simulate a pending spec change: bump generation + session.getMetadata().setGeneration(2L); + return Optional.of(session); + } + + @Override + public KubernetesClient getClient() { + return kubernetesClient; + } + }; + + reconciler.reconcile(sessionJob, mismatchContext); + assertEquals(0, flinkService.listJobs().size()); + } + + @Test + public void testSessionClusterReadyWhenRolledBack() throws Exception { + FlinkSessionJob sessionJob = TestUtils.buildSessionJob(); + + // Create a context where JM=READY, state=ROLLED_BACK, generation matches + Context rolledBackContext = + new TestUtils.TestingContext<>() { + @Override + public Optional getSecondaryResource( + Class expectedType, String eventSourceName) { + var session = TestUtils.buildSessionCluster(); + session.getStatus() + .setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); + session.getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(session.getSpec(), session); + session.getStatus().getReconciliationStatus().markReconciledSpecAsStable(); + session.getStatus() + .getReconciliationStatus() + .setState(ReconciliationState.ROLLED_BACK); + return Optional.of(session); + } + + @Override + public KubernetesClient getClient() { + return kubernetesClient; + } + }; + + reconciler.reconcile(sessionJob, rolledBackContext); + assertEquals(1, flinkService.listJobs().size()); + } + + @Test + public void testSessionClusterReadyWhenDeployed() throws Exception { + FlinkSessionJob sessionJob = TestUtils.buildSessionJob(); + + // Create a context where JM=READY, state=DEPLOYED, generation matches + Context rolledBackContext = + new TestUtils.TestingContext<>() { + @Override + public Optional getSecondaryResource( + Class expectedType, String eventSourceName) { + var session = TestUtils.buildSessionCluster(); + session.getStatus() + .setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); + session.getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(session.getSpec(), session); + session.getStatus().getReconciliationStatus().markReconciledSpecAsStable(); + session.getStatus() + .getReconciliationStatus() + .setState(ReconciliationState.DEPLOYED); + return Optional.of(session); + } + + @Override + public KubernetesClient getClient() { + return kubernetesClient; + } + }; + + reconciler.reconcile(sessionJob, rolledBackContext); + assertEquals(1, flinkService.listJobs().size()); + } } From b4755e604a5c161796c03ef6f6b8fda06d7aa5cd Mon Sep 17 00:00:00 2001 From: Ihor Mielientiev Date: Fri, 27 Feb 2026 09:44:15 +0100 Subject: [PATCH 2/3] Make sure cleanup works correctly --- .../sessionjob/SessionJobReconciler.java | 88 +++++++++++-------- .../deployment/SessionReconcilerTest.java | 25 +++++- .../sessionjob/SessionJobReconcilerTest.java | 29 +++++- 3 files changed, 101 insertions(+), 41 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java index f70a65841b..e6dfe5b2aa 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconciler.java @@ -62,8 +62,9 @@ public SessionJobReconciler( @Override public boolean readyToReconcile(FlinkResourceContext ctx) { - return sessionClusterReady( - ctx.getJosdkContext().getSecondaryResource(FlinkDeployment.class)) + var flinkDeploymentOpt = ctx.getJosdkContext().getSecondaryResource(FlinkDeployment.class); + return sessionClusterReady(flinkDeploymentOpt) + && noDeploymentChangesPending(flinkDeploymentOpt) && super.readyToReconcile(ctx); } @@ -180,44 +181,59 @@ public DeleteControl cleanupInternal(FlinkResourceContext ctx) } public static boolean sessionClusterReady(Optional flinkDeploymentOpt) { - if (flinkDeploymentOpt.isPresent()) { - var flinkdep = flinkDeploymentOpt.get(); - var jobmanagerDeploymentStatus = flinkdep.getStatus().getJobManagerDeploymentStatus(); - if (jobmanagerDeploymentStatus != JobManagerDeploymentStatus.READY) { - LOG.info( - "Session cluster deployment is in {} status, not ready for serve", - jobmanagerDeploymentStatus); - return false; - } + if (flinkDeploymentOpt.isEmpty()) { + LOG.warn("Session cluster deployment is not found"); + return false; + } + var flinkdep = flinkDeploymentOpt.get(); + var jobmanagerDeploymentStatus = flinkdep.getStatus().getJobManagerDeploymentStatus(); + if (jobmanagerDeploymentStatus != JobManagerDeploymentStatus.READY) { + LOG.info( + "Session cluster deployment is in {} status, not ready for serve", + jobmanagerDeploymentStatus); + return false; + } - // Block while FlinkDeployment is in reconciliation state - var reconciliationState = flinkdep.getStatus().getReconciliationStatus().getState(); - if (reconciliationState != ReconciliationState.DEPLOYED - && reconciliationState != ReconciliationState.ROLLED_BACK) { - LOG.info( - "Session cluster deployment reconciliation state is {}, not ready for sync", - reconciliationState); - return false; - } + // Block while FlinkDeployment is in a transitional reconciliation state + var reconciliationState = flinkdep.getStatus().getReconciliationStatus().getState(); + if (reconciliationState != ReconciliationState.DEPLOYED + && reconciliationState != ReconciliationState.ROLLED_BACK) { + LOG.info( + "Session cluster deployment reconciliation state is {}, not ready for serve", + reconciliationState); + return false; + } - // Block while a FlinkDeployment spec change is pending. This make sure that - // FlinkSessionJob updates will be - // applied only after all changes for this generation got applied - if (!Objects.equals( - flinkdep.getMetadata().getGeneration(), - flinkdep.getStatus().getObservedGeneration())) { - LOG.info( - "Session cluster deployment has pending spec changes " - + "(generation={}, observedGeneration={}), not ready for sync", - flinkdep.getMetadata().getGeneration(), - flinkdep.getStatus().getObservedGeneration()); - return false; - } + return true; + } - return true; - } else { - LOG.warn("Session cluster deployment is not found"); + /** + * Checks that the FlinkDeployment has no pending spec changes + * + *

When they differ, the deployment spec was changed but the operator hasn't acted yet. The + * cluster may still appear healthy (JM READY, state DEPLOYED), but it is about to be updated + * (e.g. deleted and recreated). Allowing a session job to start upgrading in this window would + * risk the savepoint being destroyed during the cluster rebuild. + * + *

This check is only used in {@link #readyToReconcile}, not in {@link #sessionClusterReady}, + * so it does not block cleanup or FlinkService creation. + */ + private static boolean noDeploymentChangesPending( + Optional flinkDeploymentOpt) { + if (flinkDeploymentOpt.isEmpty()) { + return false; + } + var flinkdep = flinkDeploymentOpt.get(); + if (!Objects.equals( + flinkdep.getMetadata().getGeneration(), + flinkdep.getStatus().getObservedGeneration())) { + LOG.info( + "Session cluster deployment has pending spec changes " + + "(generation={}, observedGeneration={}), not ready for sync", + flinkdep.getMetadata().getGeneration(), + flinkdep.getStatus().getObservedGeneration()); return false; } + return true; } } diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java index 37028bdd6c..46e9c3c533 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconcilerTest.java @@ -314,7 +314,6 @@ public void testUpgradeDeferredWhenSessionJobUpgrading() throws Exception { ReconciliationState.DEPLOYED, deployment.getStatus().getReconciliationStatus().getState()); - // Create a session job that is in UPGRADING state FlinkSessionJob upgradingJob = TestUtils.buildSessionJob(); upgradingJob .getStatus() @@ -323,7 +322,6 @@ public void testUpgradeDeferredWhenSessionJobUpgrading() throws Exception { upgradingJob.getStatus().getReconciliationStatus().setState(ReconciliationState.UPGRADING); kubernetesClient.resource(upgradingJob).createOrReplace(); - // Simulate a pending spec change on the deployment deployment.getMetadata().setGeneration(2L); deployment.getSpec().setRestartNonce(2L); @@ -340,13 +338,11 @@ public void testUpgradeDeferredWhenSessionJobUpgrading() throws Exception { public void testUpgradeProceedsWhenSessionJobDeployed() throws Exception { FlinkDeployment deployment = TestUtils.buildSessionCluster(); - // Initial deploy reconciler.reconcile(deployment, flinkService.getContext()); assertEquals( ReconciliationState.DEPLOYED, deployment.getStatus().getReconciliationStatus().getState()); - // Create a session job that is in DEPLOYED state (not upgrading) FlinkSessionJob deployedJob = TestUtils.buildSessionJob(); deployedJob .getStatus() @@ -374,4 +370,25 @@ public void testUpgradeProceedsWhenSessionJobDeployed() throws Exception { .deserializeLastReconciledSpec() .getRestartNonce()); } + + @Test + public void testDeploymentCleanupWaitsForSessionJobs() throws Exception { + FlinkDeployment deployment = TestUtils.buildSessionCluster(); + reconciler.reconcile(deployment, flinkService.getContext()); + + FlinkSessionJob job = TestUtils.buildSessionJob(); + job.getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(job.getSpec(), job); + job.getStatus().getReconciliationStatus().setState(ReconciliationState.DEPLOYED); + kubernetesClient.resource(job).createOrReplace(); + + var deleteControl = reconciler.cleanup(deployment, flinkService.getContext()); + assertFalse(deleteControl.isRemoveFinalizer()); + + kubernetesClient.resource(job).delete(); + + deleteControl = reconciler.cleanup(deployment, flinkService.getContext()); + assertTrue(deleteControl.isRemoveFinalizer()); + } } diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java index e0370a4198..8754ea5c60 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/sessionjob/SessionJobReconcilerTest.java @@ -882,7 +882,6 @@ public Optional getSecondaryResource( session.getStatus() .getReconciliationStatus() .serializeAndSetLastReconciledSpec(session.getSpec(), session); - // Leave state as UPGRADING (the default) return Optional.of(session); } @@ -995,4 +994,32 @@ public KubernetesClient getClient() { reconciler.reconcile(sessionJob, rolledBackContext); assertEquals(1, flinkService.listJobs().size()); } + + @Test + public void testSessionJobCleanupWorksWhenDeploymentReady() throws Exception { + FlinkSessionJob sessionJob = TestUtils.buildSessionJob(); + var readyContext = TestUtils.createContextWithReadyFlinkDeployment(kubernetesClient); + + // Deploy session job and set it to running + reconciler.reconcile(sessionJob, readyContext); + assertEquals(1, flinkService.listJobs().size()); + sessionJob + .getStatus() + .getJobStatus() + .setState(org.apache.flink.api.common.JobStatus.RUNNING); + + // First cleanup starts cancellation — defers to wait for it + var deleteControl = reconciler.cleanup(sessionJob, readyContext); + assertFalse(deleteControl.isRemoveFinalizer()); + + // Job reaches cancelled state + sessionJob + .getStatus() + .getJobStatus() + .setState(org.apache.flink.api.common.JobStatus.CANCELED); + + // Second cleanup sees terminal state — removes finalizer + deleteControl = reconciler.cleanup(sessionJob, readyContext); + assertTrue(deleteControl.isRemoveFinalizer()); + } } From 2a0584e04e6a3c66c138a8c843a703e61cde8cef Mon Sep 17 00:00:00 2001 From: Ihor Mielientiev Date: Fri, 27 Feb 2026 10:49:31 +0100 Subject: [PATCH 3/3] improve readability --- .../reconciler/deployment/SessionReconciler.java | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java index bc86f37733..c7ff0b44da 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/SessionReconciler.java @@ -71,16 +71,12 @@ protected boolean readyToReconcile(FlinkResourceContext ctx) { if (!Objects.equals( deployment.getMetadata().getGeneration(), deployment.getStatus().getObservedGeneration())) { - boolean hasUpgradingSessionJob = + var deploymentName = deployment.getMetadata().getName(); + boolean hasTransitioningSessionJob = ctx.getJosdkContext().getSecondaryResources(FlinkSessionJob.class).stream() - .filter( - job -> - deployment - .getMetadata() - .getName() - .equals(job.getSpec().getDeploymentName())) + .filter(job -> deploymentName.equals(job.getSpec().getDeploymentName())) .anyMatch(SessionReconciler::isSessionJobTransitioning); - if (hasUpgradingSessionJob) { + if (hasTransitioningSessionJob) { LOG.info( "Deferring session cluster upgrade: associated session jobs are being upgraded"); return false;