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..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 @@ -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,38 @@ 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())) { + var deploymentName = deployment.getMetadata().getName(); + boolean hasTransitioningSessionJob = + ctx.getJosdkContext().getSecondaryResources(FlinkSessionJob.class).stream() + .filter(job -> deploymentName.equals(job.getSpec().getDeploymentName())) + .anyMatch(SessionReconciler::isSessionJobTransitioning); + if (hasTransitioningSessionJob) { + 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..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 @@ -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}. */ @@ -60,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); } @@ -178,20 +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; - } else { - return true; - } - } else { + 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 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; + } + + return true; + } + + /** + * 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/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..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 @@ -305,4 +305,90 @@ 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()); + + FlinkSessionJob upgradingJob = TestUtils.buildSessionJob(); + upgradingJob + .getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(upgradingJob.getSpec(), upgradingJob); + upgradingJob.getStatus().getReconciliationStatus().setState(ReconciliationState.UPGRADING); + kubernetesClient.resource(upgradingJob).createOrReplace(); + + 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(); + + reconciler.reconcile(deployment, flinkService.getContext()); + assertEquals( + ReconciliationState.DEPLOYED, + deployment.getStatus().getReconciliationStatus().getState()); + + 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()); + } + + @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 b001ed7739..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 @@ -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,161 @@ 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); + 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()); + } + + @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()); + } }