From 90edc9c6586fc451af5ef6ecb0c3c6b5d029f588 Mon Sep 17 00:00:00 2001 From: Gyula Fora Date: Mon, 5 Jan 2026 14:05:08 +0100 Subject: [PATCH 1/8] [FLINK-38859] Do not apply minimum job vertex parallelism for whole job in native mode --- .../operator/config/FlinkConfigBuilder.java | 8 +++-- .../config/FlinkConfigBuilderTest.java | 32 +++++++++++++++++++ 2 files changed, 37 insertions(+), 3 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilder.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilder.java index e50e1268cc..b99863c809 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilder.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilder.java @@ -363,9 +363,11 @@ private int getParallelism() { * effectiveConfig.get(TaskManagerOptions.NUM_TASK_SLOTS); } - Optional maxOverrideParallelism = getMaxParallelismFromOverrideConfig(); - if (maxOverrideParallelism.isPresent() && maxOverrideParallelism.get() > 0) { - return maxOverrideParallelism.get(); + if (KubernetesDeploymentMode.STANDALONE.equals(spec.getMode())) { + Optional maxOverrideParallelism = getMaxParallelismFromOverrideConfig(); + if (maxOverrideParallelism.isPresent() && maxOverrideParallelism.get() > 0) { + return maxOverrideParallelism.get(); + } } return spec.getJob().getParallelism(); diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilderTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilderTest.java index 674ef4d578..0d2c2c93fd 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilderTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/config/FlinkConfigBuilderTest.java @@ -874,6 +874,38 @@ public void testApplyStandaloneSessionSpec() throws URISyntaxException, IOExcept StandaloneKubernetesConfigOptionsInternal.KUBERNETES_TASKMANAGER_REPLICAS)); } + @Test + public void testParallelismOverridesOnlyAppliedForStandaloneMode() + throws URISyntaxException, IOException { + FlinkDeployment dep = ReconciliationUtils.clone(flinkDeployment); + dep.getSpec().setTaskManager(new TaskManagerSpec()); + dep.getSpec().getJob().setParallelism(5); + dep.getSpec() + .getFlinkConfiguration() + .put(PipelineOptions.PARALLELISM_OVERRIDES.key(), "vertex1:10,vertex2:20"); + + // Test STANDALONE mode - parallelism overrides should be used + dep.getSpec().setMode(KubernetesDeploymentMode.STANDALONE); + Configuration configuration = + new FlinkConfigBuilder(dep, new Configuration()) + .applyFlinkConfiguration() + .applyTaskManagerSpec() + .applyJobOrSessionSpec() + .build(); + assertEquals(20, configuration.get(CoreOptions.DEFAULT_PARALLELISM)); + + // Test NATIVE mode - parallelism overrides should NOT be used, fall back to job + // parallelism + dep.getSpec().setMode(KubernetesDeploymentMode.NATIVE); + configuration = + new FlinkConfigBuilder(dep, new Configuration()) + .applyFlinkConfiguration() + .applyTaskManagerSpec() + .applyJobOrSessionSpec() + .build(); + assertEquals(5, configuration.get(CoreOptions.DEFAULT_PARALLELISM)); + } + @Test public void testBuildFrom() throws Exception { final Configuration configuration = From e480eb724ef0a921a759af2a0a87cb77c21187bf Mon Sep 17 00:00:00 2001 From: Gyula Fora Date: Mon, 5 Jan 2026 14:49:08 +0100 Subject: [PATCH 2/8] [FLINK-38858] Reuse jobid after failed submissions to guard against job duplication --- ...ernetes_operator_config_configuration.html | 6 +++++ .../shortcodes/generated/system_section.html | 6 +++++ .../config/FlinkOperatorConfiguration.java | 7 +++++- .../KubernetesOperatorConfigOptions.java | 7 ++++++ .../sessionjob/SessionJobReconciler.java | 24 +++++++++++++------ .../service/AbstractFlinkService.java | 2 +- .../operator/TestingFlinkService.java | 5 ++++ .../FlinkSessionJobControllerTest.java | 12 +++++++++- 8 files changed, 59 insertions(+), 10 deletions(-) diff --git a/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html b/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html index db57a55382..e92f04f5ca 100644 --- a/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html +++ b/docs/layouts/shortcodes/generated/kubernetes_operator_config_configuration.html @@ -212,6 +212,12 @@ Boolean Indicate whether a savepoint must be taken when deleting a FlinkDeployment or FlinkSessionJob. + +
kubernetes.operator.job.submission.timeout
+ 10 min + Duration + The timeout for session job submissions. +
kubernetes.operator.job.upgrade.ignore-pending-savepoint
false diff --git a/docs/layouts/shortcodes/generated/system_section.html b/docs/layouts/shortcodes/generated/system_section.html index 6331b62a0c..6add6cc021 100644 --- a/docs/layouts/shortcodes/generated/system_section.html +++ b/docs/layouts/shortcodes/generated/system_section.html @@ -50,6 +50,12 @@ Duration The timeout for the observer to wait the flink rest client to return. + +
kubernetes.operator.job.submission.timeout
+ 10 min + Duration + The timeout for session job submissions. +
kubernetes.operator.leader-election.enabled
false diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java index 97068722aa..0733328214 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkOperatorConfiguration.java @@ -80,6 +80,7 @@ public class FlinkOperatorConfiguration { int reportedExceptionEventsMaxCount; int reportedExceptionEventsMaxStackTraceLength; boolean manageIngress; + Duration jobSubmissionTimeout; public static FlinkOperatorConfiguration fromConfiguration(Configuration operatorConfig) { Duration reconcileInterval = @@ -207,6 +208,9 @@ public static FlinkOperatorConfiguration fromConfiguration(Configuration operato boolean manageIngress = operatorConfig.get(KubernetesOperatorConfigOptions.OPERATOR_MANAGE_INGRESS); + Duration jobSubmissionTimeout = + operatorConfig.get(KubernetesOperatorConfigOptions.OPERATOR_JOB_SUBMISSION_TIMEOUT); + return new FlinkOperatorConfiguration( reconcileInterval, reconcilerMaxParallelism, @@ -239,7 +243,8 @@ public static FlinkOperatorConfiguration fromConfiguration(Configuration operato slowRequestThreshold, reportedExceptionEventsMaxCount, reportedExceptionEventsMaxStackTraceLength, - manageIngress); + manageIngress, + jobSubmissionTimeout); } private static GenericRetry getRetryConfig(Configuration conf) { diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java index 4f755372a1..f0dc943030 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/KubernetesOperatorConfigOptions.java @@ -134,6 +134,13 @@ public static String operatorConfigKey(String key) { .withDescription( "The timeout for the resource clean up to wait for flink to shutdown cluster."); + @Documentation.Section(SECTION_SYSTEM) + public static final ConfigOption OPERATOR_JOB_SUBMISSION_TIMEOUT = + operatorConfig("job.submission.timeout") + .durationType() + .defaultValue(Duration.ofMinutes(10)) + .withDescription("The timeout for session job submissions."); + @Documentation.Section(SECTION_DYNAMIC) public static final ConfigOption DEPLOYMENT_ROLLBACK_ENABLED = operatorConfig("deployment.rollback.enabled") 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 4932940a4f..f5d10626b1 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 @@ -18,6 +18,7 @@ package org.apache.flink.kubernetes.operator.reconciler.sessionjob; import org.apache.flink.api.common.JobID; +import org.apache.flink.api.common.JobStatus; import org.apache.flink.autoscaler.JobAutoScaler; import org.apache.flink.configuration.Configuration; import org.apache.flink.kubernetes.operator.api.FlinkDeployment; @@ -81,10 +82,22 @@ public void deploy( MSG_SUBMIT, ctx.getKubernetesClient()); - // Generate job id and record in status for durability - var jobId = JobID.generate(); - ctx.getResource().getStatus().getJobStatus().setJobId(jobId.toHexString()); - statusRecorder.patchAndCacheStatus(ctx.getResource(), ctx.getKubernetesClient()); + var jobStatus = ctx.getResource().getStatus().getJobStatus(); + + String existingJobIdStr = jobStatus.getJobId(); + JobID jobId; + var jobState = jobStatus.getState(); + var reuseJobId = jobState == JobStatus.RECONCILING; + if (existingJobIdStr != null && reuseJobId) { + jobId = JobID.fromHexString(existingJobIdStr); + LOG.info("Reusing existing job ID {} for deployment retry", jobId); + } else { + jobId = JobID.generate(); + LOG.info("Generated new job ID {} for deployment", jobId); + jobStatus.setJobId(jobId.toHexString()); + jobStatus.setState(org.apache.flink.api.common.JobStatus.RECONCILING); + statusRecorder.patchAndCacheStatus(ctx.getResource(), ctx.getKubernetesClient()); + } ctx.getFlinkService() .submitJobToSessionCluster( @@ -93,9 +106,6 @@ public void deploy( jobId, deployConfig, savepoint.orElse(null)); - - var status = ctx.getResource().getStatus(); - status.getJobStatus().setState(org.apache.flink.api.common.JobStatus.RECONCILING); } @Override diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java index 820fbced76..5598410c41 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java @@ -908,7 +908,7 @@ protected void runJar( LOG.info("Submitting job: {} to session cluster.", jobID); clusterClient .sendRequest(headers, parameters, runRequestBody) - .get(operatorConfig.getFlinkClientTimeout().toSeconds(), TimeUnit.SECONDS); + .get(operatorConfig.getJobSubmissionTimeout().toSeconds(), TimeUnit.SECONDS); } catch (Exception e) { LOG.error("Failed to submit job to session cluster.", e); throw new FlinkRuntimeException(e); diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java index 8440ffe8e7..b17fb624a1 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java @@ -51,6 +51,7 @@ import org.apache.flink.kubernetes.operator.service.CheckpointHistoryWrapper; import org.apache.flink.kubernetes.operator.service.SuspendMode; import org.apache.flink.kubernetes.operator.standalone.StandaloneKubernetesConfigOptionsInternal; +import org.apache.flink.runtime.client.DuplicateJobSubmissionException; import org.apache.flink.runtime.client.JobStatusMessage; import org.apache.flink.runtime.clusterframework.ApplicationStatus; import org.apache.flink.runtime.execution.ExecutionState; @@ -292,6 +293,10 @@ public JobID submitJobToSessionCluster( if (deployFailure) { throw new Exception("Deployment failure"); } + if (jobs.stream().anyMatch(job -> job.f1.getJobId().equals(jobID))) { + throw DuplicateJobSubmissionException.of(jobID); + } + JobStatusMessage jobStatusMessage = new JobStatusMessage( jobID, diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java index 09ff806e8c..7dee2c8317 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java @@ -88,7 +88,7 @@ public void before() { } @Test - public void testSubmitJobButException() { + public void testSubmitJobButException() throws Exception { flinkService.setDeployFailure(true); try { @@ -97,6 +97,10 @@ public void testSubmitJobButException() { // Ignore } + assertEquals(sessionJob.getStatus().getJobStatus().getState(), RECONCILING); + String jobId = sessionJob.getStatus().getJobStatus().getJobId(); + assertNotNull(jobId); + Assertions.assertEquals(2, testController.events().size()); // Discard submit event testController.events().remove(); @@ -105,6 +109,12 @@ public void testSubmitJobButException() { Assertions.assertEquals(EventRecorder.Type.Warning.toString(), event.getType()); Assertions.assertEquals("Error", event.getReason()); + flinkService.setDeployFailure(false); + testController.reconcile(sessionJob, context); + + // Make sure we reused the original failed job id + assertEquals(jobId, flinkService.listJobs().get(0).f1.getJobId().toHexString()); + testController.cleanup(sessionJob, context); } From 723eb75c17e3a001e6d425c2474414c59b8baf16 Mon Sep 17 00:00:00 2001 From: Gyula Fora Date: Mon, 5 Jan 2026 14:49:32 +0100 Subject: [PATCH 3/8] [FLINK-38860] Remove flink shaded guava usage --- flink-autoscaler-plugin-jdbc/pom.xml | 5 +++++ .../jdbc/event/JdbcAutoScalerEventHandler.java | 3 +-- flink-autoscaler-standalone/pom.xml | 5 +++++ .../standalone/StandaloneAutoscalerExecutor.java | 3 +-- flink-autoscaler/pom.xml | 5 +++++ .../apache/flink/autoscaler/RestApiMetricsCollector.java | 3 +-- .../apache/flink/autoscaler/topology/JobTopology.java | 4 ++-- flink-kubernetes-operator/pom.xml | 5 +++++ .../kubeclient/decorators/FlinkConfMountDecorator.java | 2 +- .../kubernetes/operator/config/FlinkConfigManager.java | 9 ++++----- .../operator/exception/DeploymentFailedException.java | 3 +-- .../operator/service/AbstractFlinkService.java | 2 +- .../operator/service/FlinkResourceContextFactory.java | 3 +-- .../src/main/resources/META-INF/NOTICE | 6 ++++++ pom.xml | 6 ++++++ tools/maven/checkstyle.xml | 5 +++-- 16 files changed, 48 insertions(+), 21 deletions(-) diff --git a/flink-autoscaler-plugin-jdbc/pom.xml b/flink-autoscaler-plugin-jdbc/pom.xml index 6add87948f..2f1c0facc5 100644 --- a/flink-autoscaler-plugin-jdbc/pom.xml +++ b/flink-autoscaler-plugin-jdbc/pom.xml @@ -73,6 +73,11 @@ under the License. provided + + com.google.guava + guava + + diff --git a/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/event/JdbcAutoScalerEventHandler.java b/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/event/JdbcAutoScalerEventHandler.java index 8faef51f48..69639a8db2 100644 --- a/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/event/JdbcAutoScalerEventHandler.java +++ b/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/event/JdbcAutoScalerEventHandler.java @@ -25,8 +25,7 @@ import org.apache.flink.runtime.jobgraph.JobVertexID; import org.apache.flink.util.Preconditions; -import org.apache.flink.shaded.guava31.com.google.common.util.concurrent.ThreadFactoryBuilder; - +import com.google.common.util.concurrent.ThreadFactoryBuilder; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; diff --git a/flink-autoscaler-standalone/pom.xml b/flink-autoscaler-standalone/pom.xml index b270863253..83b5c1cb22 100644 --- a/flink-autoscaler-standalone/pom.xml +++ b/flink-autoscaler-standalone/pom.xml @@ -45,6 +45,11 @@ under the License. ${project.version} + + com.google.guava + guava + + org.apache.flink flink-runtime diff --git a/flink-autoscaler-standalone/src/main/java/org/apache/flink/autoscaler/standalone/StandaloneAutoscalerExecutor.java b/flink-autoscaler-standalone/src/main/java/org/apache/flink/autoscaler/standalone/StandaloneAutoscalerExecutor.java index 7e13be9288..c5de149cdd 100644 --- a/flink-autoscaler-standalone/src/main/java/org/apache/flink/autoscaler/standalone/StandaloneAutoscalerExecutor.java +++ b/flink-autoscaler-standalone/src/main/java/org/apache/flink/autoscaler/standalone/StandaloneAutoscalerExecutor.java @@ -26,8 +26,7 @@ import org.apache.flink.configuration.UnmodifiableConfiguration; import org.apache.flink.util.concurrent.ExecutorThreadFactory; -import org.apache.flink.shaded.guava31.com.google.common.util.concurrent.ThreadFactoryBuilder; - +import com.google.common.util.concurrent.ThreadFactoryBuilder; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.slf4j.MDC; diff --git a/flink-autoscaler/pom.xml b/flink-autoscaler/pom.xml index de63218b1b..d97f33555d 100644 --- a/flink-autoscaler/pom.xml +++ b/flink-autoscaler/pom.xml @@ -78,6 +78,11 @@ under the License. provided + + com.google.guava + guava + + org.apache.flink flink-test-utils-junit diff --git a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/RestApiMetricsCollector.java b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/RestApiMetricsCollector.java index 67d77c3523..7ebedb76ba 100644 --- a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/RestApiMetricsCollector.java +++ b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/RestApiMetricsCollector.java @@ -35,8 +35,7 @@ import org.apache.flink.runtime.rest.messages.job.metrics.MetricsAggregationParameter; import org.apache.flink.runtime.rest.messages.job.metrics.MetricsFilterParameter; -import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableMap; - +import com.google.common.collect.ImmutableMap; import lombok.SneakyThrows; import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; diff --git a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/topology/JobTopology.java b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/topology/JobTopology.java index 83e8111701..2600450941 100644 --- a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/topology/JobTopology.java +++ b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/topology/JobTopology.java @@ -21,14 +21,14 @@ import org.apache.flink.runtime.instance.SlotSharingGroupId; import org.apache.flink.runtime.jobgraph.JobVertexID; -import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableMap; -import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableSet; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ArrayNode; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.ImmutableSet; import lombok.EqualsAndHashCode; import lombok.Getter; import lombok.ToString; diff --git a/flink-kubernetes-operator/pom.xml b/flink-kubernetes-operator/pom.xml index 618ffa84d9..728ef11cd2 100644 --- a/flink-kubernetes-operator/pom.xml +++ b/flink-kubernetes-operator/pom.xml @@ -154,6 +154,11 @@ under the License. ${log4j.version} + + com.google.guava + guava + + org.apache.flink diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java index 4615423ba8..2d5cbe6d9e 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java @@ -43,7 +43,7 @@ import org.apache.flink.kubernetes.shaded.io.fabric8.kubernetes.api.model.VolumeBuilder; import org.apache.flink.kubernetes.utils.Constants; -import org.apache.flink.shaded.guava31.com.google.common.io.Files; +import com.google.common.io.Files; import java.io.File; import java.io.IOException; diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigManager.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigManager.java index f31fc7ba4b..62fe3239c7 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigManager.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/config/FlinkConfigManager.java @@ -33,13 +33,12 @@ import org.apache.flink.kubernetes.operator.utils.EnvUtils; import org.apache.flink.kubernetes.operator.utils.FlinkUtils; -import org.apache.flink.shaded.guava31.com.google.common.cache.Cache; -import org.apache.flink.shaded.guava31.com.google.common.cache.CacheBuilder; -import org.apache.flink.shaded.guava31.com.google.common.cache.CacheLoader; -import org.apache.flink.shaded.guava31.com.google.common.cache.LoadingCache; - import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.cache.Cache; +import com.google.common.cache.CacheBuilder; +import com.google.common.cache.CacheLoader; +import com.google.common.cache.LoadingCache; import io.fabric8.kubernetes.api.model.ObjectMeta; import lombok.Builder; import lombok.SneakyThrows; diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/DeploymentFailedException.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/DeploymentFailedException.java index 9fc6143f06..7ce84243a2 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/DeploymentFailedException.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/exception/DeploymentFailedException.java @@ -17,8 +17,7 @@ package org.apache.flink.kubernetes.operator.exception; -import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableSet; - +import com.google.common.collect.ImmutableSet; import io.fabric8.kubernetes.api.model.ContainerStatus; import io.fabric8.kubernetes.api.model.apps.DeploymentCondition; diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java index 5598410c41..f0a7b25bb5 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java @@ -108,9 +108,9 @@ import org.apache.flink.util.FlinkRuntimeException; import org.apache.flink.util.Preconditions; -import org.apache.flink.shaded.guava31.com.google.common.collect.Iterables; import org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpResponseStatus; +import com.google.common.collect.Iterables; import io.fabric8.kubernetes.api.model.DeletionPropagation; import io.fabric8.kubernetes.api.model.ObjectMeta; import io.fabric8.kubernetes.api.model.PodList; diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkResourceContextFactory.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkResourceContextFactory.java index f2f75a5295..9cc9b489cf 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkResourceContextFactory.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkResourceContextFactory.java @@ -35,8 +35,7 @@ import org.apache.flink.kubernetes.operator.utils.EventRecorder; import org.apache.flink.util.concurrent.ExecutorThreadFactory; -import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableMap; - +import com.google.common.collect.ImmutableMap; import io.javaoperatorsdk.operator.api.reconciler.Context; import io.javaoperatorsdk.operator.processing.event.ResourceID; import lombok.Data; diff --git a/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE b/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE index b63fb8801e..497b3c7c06 100644 --- a/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE +++ b/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE @@ -12,6 +12,11 @@ This project bundles the following dependencies under the Apache Software Licens - com.fasterxml.jackson.dataformat:jackson-dataformat-yaml:jar:2.15.0 - com.fasterxml.jackson.datatype:jackson-datatype-jsr310:jar:2.15.0 - com.google.code.findbugs:jsr305:jar:1.3.9 +- com.google.errorprone:error_prone_annotations:jar:2.36.0 +- com.google.guava:failureaccess:jar:1.0.2 +- com.google.guava:guava:jar:33.4.0-jre +- com.google.guava:listenablefuture:jar:9999.0-empty-to-avoid-conflict-with-guava +- com.google.j2objc:j2objc-annotations:jar:3.0.0 - com.squareup.okhttp3:logging-interceptor:jar:4.12.0 - com.squareup.okhttp3:okhttp:jar:4.12.0 - com.squareup.okio:okio-jvm:jar:3.6.0 @@ -58,6 +63,7 @@ This project bundles the following dependencies under the Apache Software Licens - org.apache.logging.log4j:log4j-api:jar:2.23.1 - org.apache.logging.log4j:log4j-core:jar:2.23.1 - org.apache.logging.log4j:log4j-slf4j-impl:jar:2.23.1 +- org.checkerframework:checker-qual:jar:3.43.0 - org.javassist:javassist:jar:3.24.0-GA - org.jetbrains.kotlin:kotlin-stdlib-common:jar:1.8.21 - org.jetbrains.kotlin:kotlin-stdlib-jdk7:jar:1.8.21 diff --git a/pom.xml b/pom.xml index 44a1fffd26..19d27b88ec 100644 --- a/pom.xml +++ b/pom.xml @@ -84,6 +84,7 @@ under the License. 3.18.0 2.17.0 1.20.1 + 33.4.0-jre 1.7.36 2.23.1 @@ -133,6 +134,11 @@ under the License. pom import + + com.google.guava + guava + ${guava.version} + io.fabric8 kubernetes-client diff --git a/tools/maven/checkstyle.xml b/tools/maven/checkstyle.xml index b160f472db..216ad4799b 100644 --- a/tools/maven/checkstyle.xml +++ b/tools/maven/checkstyle.xml @@ -230,8 +230,9 @@ This file is based on the checkstyle file of Apache Beam. - - + + + From 824d96facaafd61d15b92ae1bf80ab8e76b80a85 Mon Sep 17 00:00:00 2001 From: Gyula Fora Date: Thu, 10 Jul 2025 12:46:06 +0200 Subject: [PATCH 4/8] [FLINK-38077] Make sure jobmanager is accessible when trying to cancel for suspend/upgrade --- .../operator/api/status/CommonStatus.java | 12 +++++ .../api/status/FlinkDeploymentStatus.java | 7 +++ .../status/JobManagerDeploymentStatus.java | 4 ++ .../reconciler/ReconciliationUtils.java | 4 -- .../deployment/AbstractJobReconciler.java | 4 +- .../deployment/ApplicationReconciler.java | 7 +-- .../service/AbstractFlinkService.java | 2 +- .../ApplicationReconcilerUpgradeModeTest.java | 45 ++++++++++++++----- 8 files changed, 63 insertions(+), 22 deletions(-) diff --git a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/CommonStatus.java b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/CommonStatus.java index 30649e8968..0de792da5a 100644 --- a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/CommonStatus.java +++ b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/CommonStatus.java @@ -22,6 +22,7 @@ import org.apache.flink.kubernetes.operator.api.spec.AbstractFlinkSpec; import org.apache.flink.kubernetes.operator.api.spec.JobState; +import com.fasterxml.jackson.annotation.JsonIgnore; import io.fabric8.crd.generator.annotation.PrinterColumn; import lombok.AllArgsConstructor; import lombok.Data; @@ -29,6 +30,8 @@ import lombok.experimental.SuperBuilder; import org.apache.commons.lang3.StringUtils; +import static org.apache.flink.api.common.JobStatus.RECONCILING; + /** Last observed common status of the Flink deployment/Flink SessionJob. */ @Experimental @Data @@ -121,4 +124,13 @@ public ResourceLifecycleState getLifecycleState() { return ResourceLifecycleState.DEPLOYED; } + + @JsonIgnore + public boolean isJobCancellable() { + var jobState = jobStatus.getState(); + if (jobState == null) { + return false; + } + return RECONCILING != jobState && !jobState.isGloballyTerminalState(); + } } diff --git a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkDeploymentStatus.java b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkDeploymentStatus.java index 136d3415f3..9a34129963 100644 --- a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkDeploymentStatus.java +++ b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/FlinkDeploymentStatus.java @@ -20,6 +20,7 @@ import org.apache.flink.annotation.Experimental; import org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec; +import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import lombok.AllArgsConstructor; import lombok.Data; @@ -55,4 +56,10 @@ public class FlinkDeploymentStatus extends CommonStatus { /** Information about the TaskManagers for the scale subresource. */ private TaskManagerInfo taskManager; + + @JsonIgnore + @Override + public boolean isJobCancellable() { + return super.isJobCancellable() && jobManagerDeploymentStatus.isRestApiAvailable(); + } } diff --git a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/JobManagerDeploymentStatus.java b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/JobManagerDeploymentStatus.java index 54a0181bc0..2c86c49ac9 100644 --- a/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/JobManagerDeploymentStatus.java +++ b/flink-kubernetes-operator-api/src/main/java/org/apache/flink/kubernetes/operator/api/status/JobManagerDeploymentStatus.java @@ -35,4 +35,8 @@ public enum JobManagerDeploymentStatus { /** Deployment in terminal error, requires spec change for reconciliation to continue. */ ERROR; + + public boolean isRestApiAvailable() { + return this == READY; + } } diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/ReconciliationUtils.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/ReconciliationUtils.java index 2936501a88..94140c3ff0 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/ReconciliationUtils.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/ReconciliationUtils.java @@ -379,10 +379,6 @@ public static boolean isJobCancelled(CommonStatus status) { return CANCELED == status.getJobStatus().getState(); } - public static boolean isJobCancellable(CommonStatus status) { - return RECONCILING != status.getJobStatus().getState(); - } - public static boolean isJobCancelling(CommonStatus status) { return status.getJobStatus() != null && CANCELLING == status.getJobStatus().getState(); } diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java index 3fcf8e46cd..9a4480d0ec 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java @@ -277,7 +277,7 @@ protected JobUpgrade getJobUpgrade(FlinkResourceContext ctx, Configuration d if (running && savepointPossible) { LOG.info("Using savepoint to upgrade Flink version"); return JobUpgrade.savepoint(false); - } else if (ReconciliationUtils.isJobCancellable(resource.getStatus())) { + } else if (resource.getStatus().isJobCancellable()) { LOG.info("Using last-state upgrade with cancellation to upgrade Flink version"); return JobUpgrade.lastStateUsingCancel(); } else { @@ -354,7 +354,7 @@ protected JobUpgrade getUpgradeModeBasedOnStateAge( private boolean allowLastStateCancel(FlinkResourceContext ctx) { var resource = ctx.getResource(); - if (!ReconciliationUtils.isJobCancellable(resource.getStatus())) { + if (!resource.getStatus().isJobCancellable()) { return false; } if (resource instanceof FlinkSessionJob) { diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java index 180fb128c9..ce8baca5e0 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java @@ -88,12 +88,13 @@ protected JobUpgrade getJobUpgrade( var status = deployment.getStatus(); var availableUpgradeMode = super.getJobUpgrade(ctx, deployConfig); - if (availableUpgradeMode.isAvailable() || !availableUpgradeMode.isAllowFallback()) { + if (availableUpgradeMode.isAvailable()) { return availableUpgradeMode; } var flinkService = ctx.getFlinkService(); - if (HighAvailabilityMode.isHighAvailabilityModeActivated(deployConfig) + if (availableUpgradeMode.isAllowFallback() + && HighAvailabilityMode.isHighAvailabilityModeActivated(deployConfig) && HighAvailabilityMode.isHighAvailabilityModeActivated(ctx.getObserveConfig()) && flinkService.isHaMetadataAvailable(deployConfig)) { LOG.info( @@ -125,7 +126,7 @@ protected JobUpgrade getJobUpgrade( "UpgradeFailed"); } - return JobUpgrade.unavailable(); + return availableUpgradeMode; } private void deleteJmThatNeverStarted( diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java index f0a7b25bb5..18fc3690da 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java @@ -336,7 +336,7 @@ protected CancelResult cancelJob( savepointPath = savepointJobOrError(clusterClient, status, conf); break; case STATELESS: - if (ReconciliationUtils.isJobCancellable(status)) { + if (status.isJobCancellable()) { try { cancelJobOrError(clusterClient, status, true); } catch (Exception ex) { diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java index 3b50af2a76..c9f705869e 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java @@ -72,6 +72,7 @@ import static org.apache.flink.api.common.JobStatus.RUNNING; import static org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.SNAPSHOT_RESOURCE_ENABLED; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNull; @@ -501,6 +502,7 @@ public void testInitialJmDeployCannotStartLegacy(UpgradeMode upgradeMode, boolea @ValueSource(booleans = {true, false}) public void testLastStateMaxCheckpointAge(boolean cancellable) throws Exception { var deployment = TestUtils.buildApplicationCluster(); + deployment.getStatus().setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); deployment .getSpec() .getFlinkConfiguration() @@ -643,6 +645,7 @@ public void testFlinkVersionSwitching( jobStatus.setJobId(new JobID().toString()); // Running state, savepoint if possible + deployment.getStatus().setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); jobStatus.setState(RUNNING); var ctx = getResourceContext(deployment); var deployConf = ctx.getDeployConfig(deployment.getSpec()); @@ -654,6 +657,7 @@ public void testFlinkVersionSwitching( jobReconciler.getJobUpgrade(ctx, deployConf)); // Not running (but cancellable) + deployment.getStatus().setJobManagerDeploymentStatus(JobManagerDeploymentStatus.READY); jobStatus.setState(RESTARTING); assertEquals( AbstractJobReconciler.JobUpgrade.lastStateUsingCancel(), @@ -667,17 +671,26 @@ public void testFlinkVersionSwitching( } private static Stream testLastStateCancelParams() { - return Stream.of( - Arguments.of(UpgradeMode.LAST_STATE, true), - Arguments.of(UpgradeMode.LAST_STATE, false), - Arguments.of(UpgradeMode.SAVEPOINT, true), - Arguments.of(UpgradeMode.SAVEPOINT, false)); + var out = new ArrayList(); + for (var upgradeMode : List.of(UpgradeMode.SAVEPOINT, UpgradeMode.LAST_STATE)) { + for (boolean allowFallback : List.of(true, false)) { + for (var jmStatus : JobManagerDeploymentStatus.values()) { + out.add(Arguments.of(upgradeMode, allowFallback, jmStatus)); + } + } + } + return out.stream(); } @ParameterizedTest @MethodSource("testLastStateCancelParams") - public void testLastStateNoHaMeta(UpgradeMode upgradeMode, boolean allowFallback) + public void testLastStateNoHaMeta( + UpgradeMode upgradeMode, boolean allowFallback, JobManagerDeploymentStatus jmStatus) throws Exception { + if (upgradeMode == UpgradeMode.LAST_STATE && !allowFallback) { + // This cannot happen + return; + } var jobReconciler = (ApplicationReconciler) this.reconciler.getReconciler(); var deployment = TestUtils.buildApplicationCluster(); deployment @@ -702,6 +715,8 @@ public void testLastStateNoHaMeta(UpgradeMode upgradeMode, boolean allowFallback // Set job status to running var jobStatus = deployment.getStatus().getJobStatus(); + deployment.getStatus().setJobManagerDeploymentStatus(jmStatus); + long now = System.currentTimeMillis(); jobStatus.setStartTime(Long.toString(now)); @@ -712,15 +727,21 @@ public void testLastStateNoHaMeta(UpgradeMode upgradeMode, boolean allowFallback var ctx = getResourceContext(deployment); var deployConf = ctx.getDeployConfig(deployment.getSpec()); - if (upgradeMode == UpgradeMode.LAST_STATE) { - assertEquals( - AbstractJobReconciler.JobUpgrade.lastStateUsingCancel(), - jobReconciler.getJobUpgrade(ctx, deployConf)); + if (List.of(JobManagerDeploymentStatus.ERROR, JobManagerDeploymentStatus.MISSING) + .contains(jmStatus)) { + assertThatThrownBy(() -> jobReconciler.getJobUpgrade(ctx, deployConf)) + .isInstanceOf(UpgradeFailureException.class) + .hasMessageContaining( + "JobManager deployment is missing and HA metadata is not available"); } else { + boolean immediatelyCancellable = + (upgradeMode == UpgradeMode.LAST_STATE || allowFallback) + && jmStatus == JobManagerDeploymentStatus.READY; assertEquals( - allowFallback + immediatelyCancellable ? AbstractJobReconciler.JobUpgrade.lastStateUsingCancel() - : AbstractJobReconciler.JobUpgrade.pendingUpgrade(), + : new AbstractJobReconciler.JobUpgrade( + null, null, false, allowFallback, true), jobReconciler.getJobUpgrade(ctx, deployConf)); } } From 57f5a7b9bfce12f9fe8600ed08f96a2a27f3da4f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Attila=20M=C3=A9sz=C3=A1ros?= Date: Wed, 14 Jan 2026 16:58:22 +0100 Subject: [PATCH 5/8] [FLINK-38875] Bump Java Operator SDK to version 5.2.2 (#1048) --- .../src/main/resources/META-INF/NOTICE | 4 ++-- .../apache/flink/kubernetes/operator/TestUtils.java | 10 ++++++++++ .../kubernetes/operator/utils/TestingJosdkContext.java | 10 ++++++++++ pom.xml | 2 +- 4 files changed, 23 insertions(+), 3 deletions(-) diff --git a/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE b/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE index 497b3c7c06..e3ed08e184 100644 --- a/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE +++ b/flink-kubernetes-operator/src/main/resources/META-INF/NOTICE @@ -53,8 +53,8 @@ This project bundles the following dependencies under the Apache Software Licens - io.fabric8:kubernetes-model-scheduling:jar:7.3.0 - io.fabric8:kubernetes-model-storageclass:jar:7.3.0 - io.fabric8:zjsonpatch:jar:7.3.0 -- io.javaoperatorsdk:operator-framework-core:jar:5.1.4 -- io.javaoperatorsdk:operator-framework:jar:5.1.4 +- io.javaoperatorsdk:operator-framework-core:jar:5.2.2 +- io.javaoperatorsdk:operator-framework:jar:5.2.2 - org.apache.commons:commons-compress:jar:1.26.0 - org.apache.commons:commons-lang3:jar:3.18.0 - org.apache.commons:commons-math3:jar:3.6.1 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 f3f2e3571d..a48f41d9c9 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 @@ -557,5 +557,15 @@ public IndexedResourceCache getPrimaryCache() { public boolean isNextReconciliationImminent() { return false; } + + @Override + public boolean isPrimaryResourceDeleted() { + return false; + } + + @Override + public boolean isPrimaryResourceFinalStateUnknown() { + return false; + } } } diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/TestingJosdkContext.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/TestingJosdkContext.java index be87f3b56d..d40ce498f9 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/TestingJosdkContext.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/TestingJosdkContext.java @@ -127,4 +127,14 @@ public IndexedResourceCache

getPrimaryCache() { public boolean isNextReconciliationImminent() { return false; } + + @Override + public boolean isPrimaryResourceDeleted() { + return false; + } + + @Override + public boolean isPrimaryResourceFinalStateUnknown() { + return false; + } } diff --git a/pom.xml b/pom.xml index 19d27b88ec..c209b27591 100644 --- a/pom.xml +++ b/pom.xml @@ -75,7 +75,7 @@ under the License. 3.3.2 5.0.0 - 5.1.4 + 5.2.2 3.0.0 7.3.1 From 9ccf27848348628fe60597fe6524c2d325b80cdd Mon Sep 17 00:00:00 2001 From: Martijn Visser <2989614+MartijnVisser@users.noreply.github.com> Date: Wed, 21 Jan 2026 11:52:17 +0100 Subject: [PATCH 6/8] [FLINK-38925][docs] Update Matomo URL to the right domain (#1054) --- .github/workflows/docs.yaml | 2 +- docs/layouts/_default/baseof.html | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/docs.yaml b/.github/workflows/docs.yaml index 20432bed8e..93f09bcf1c 100644 --- a/.github/workflows/docs.yaml +++ b/.github/workflows/docs.yaml @@ -82,7 +82,7 @@ jobs: -Dcheckstyle.skip=true \ -Dspotless.check.skip=true \ -Denforcer.skip=true \ - -Dheader="

Back to Flink Website

" + -Dheader="

Back to Flink Website

" mv target/site/apidocs docs/target/api/java - name: Upload documentation diff --git a/docs/layouts/_default/baseof.html b/docs/layouts/_default/baseof.html index 733bf1b6d2..0db6bc887b 100644 --- a/docs/layouts/_default/baseof.html +++ b/docs/layouts/_default/baseof.html @@ -34,7 +34,7 @@ _paq.push(['trackPageView']); _paq.push(['enableLinkTracking']); (function() { - var u="//matomo.privacy.apache.org/"; + var u="//analytics.apache.org/"; _paq.push(['setTrackerUrl', u+'matomo.php']); _paq.push(['setSiteId', '1']); var d=document, g=d.createElement('script'), s=d.getElementsByTagName('script')[0]; From 0b14606e15fd6182d582029cf7427a7ed40e86fd Mon Sep 17 00:00:00 2001 From: James Kan Date: Thu, 11 Dec 2025 12:15:22 -0800 Subject: [PATCH 7/8] [FLINK-38867] FlinkBlueGreenDeployment used initialSavepointPath on restartSavepointNonce Co-authored-by: Daniel Rossos --- .../api/bluegreen/BlueGreenDiffType.java | 8 +- .../bluegreen/BlueGreenDeploymentService.java | 57 ++++++- .../FlinkBlueGreenDeploymentSpecDiff.java | 21 ++- .../utils/bluegreen/BlueGreenUtils.java | 17 ++- ...linkBlueGreenDeploymentControllerTest.java | 62 ++++++++ .../FlinkBlueGreenDeploymentSpecDiffTest.java | 57 ++++++- .../utils/bluegreen/BlueGreenUtilsTest.java | 141 ++++++++++++++++++ 7 files changed, 347 insertions(+), 16 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java index 5581493f7e..e456a89133 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java @@ -26,5 +26,11 @@ public enum BlueGreenDiffType { TRANSITION, /** Changes that only affect the child FlinkDeploymentSpec. */ - PATCH_CHILD + PATCH_CHILD, + + /** + * Full redeploy from user-specified savepoint. Triggered when savepointRedeployNonce changes. + * Uses the initialSavepointPath from spec instead of taking a new savepoint. + */ + SAVEPOINT_REDEPLOY } 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 85de365b97..132b05e27f 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 @@ -42,6 +42,7 @@ import org.slf4j.LoggerFactory; import java.time.Instant; +import java.util.Objects; import static org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenKubernetesService.deleteFlinkDeployment; import static org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenKubernetesService.deployCluster; @@ -143,6 +144,25 @@ public UpdateControl checkAndInitiateDeployment( context.getDeploymentStatus().setSavepointTriggerId(null); return markDeploymentFailing(context, error); } + + } else if (specDiff == BlueGreenDiffType.SAVEPOINT_REDEPLOY) { + // Savepoint redeploy: skip taking a new savepoint, use initialSavepointPath + var jobSpec = + context.getBgDeployment().getSpec().getTemplate().getSpec().getJob(); + LOG.info( + "Savepoint redeploy triggered for '{}', using initialSavepointPath: {}", + context.getBgDeployment().getMetadata().getName(), + Objects.toString(jobSpec.getInitialSavepointPath(), "")); + setLastReconciledSpec(context); + try { + return startSavepointRedeployTransition( + context, currentBlueGreenDeploymentType); + } catch (Exception e) { + var error = + "Could not start Savepoint Redeploy Transition. Details: " + + e.getMessage(); + return markDeploymentFailing(context, error); + } } else { setLastReconciledSpec(context); LOG.info( @@ -170,6 +190,13 @@ public UpdateControl checkAndInitiateDeployment( private UpdateControl patchFlinkDeployment( BlueGreenContext context, BlueGreenDeploymentType blueGreenDeploymentTypeToPatch) { + return patchFlinkDeployment(context, blueGreenDeploymentTypeToPatch, true); + } + + private UpdateControl patchFlinkDeployment( + BlueGreenContext context, + BlueGreenDeploymentType blueGreenDeploymentTypeToPatch, + boolean carryOverSavepointInPatch) { String childDeploymentName = context.getBgDeployment().getMetadata().getName() @@ -188,7 +215,10 @@ private UpdateControl patchFlinkDeployment( // will it be used by this patching? otherwise this is unnecessary, keep lastSavepoint = // null. Savepoint lastSavepoint = - carryOverSavepoint(context, blueGreenDeploymentTypeToPatch, childDeploymentName); + carryOverSavepointInPatch + ? carryOverSavepoint( + context, blueGreenDeploymentTypeToPatch, childDeploymentName) + : null; return initiateDeployment( context, @@ -246,6 +276,26 @@ private UpdateControl startTransition( false); } + /** + * Starts a transition for savepoint redeploy scenario. Unlike normal transitions, this does not + * take a new savepoint - it uses the initialSavepointPath specified in the spec. + * + * @param context the transition context + * @param currentBlueGreenDeploymentType the current deployment type + * @return UpdateControl for the deployment + */ + private UpdateControl startSavepointRedeployTransition( + BlueGreenContext context, BlueGreenDeploymentType currentBlueGreenDeploymentType) { + DeploymentTransition transition = calculateTransition(currentBlueGreenDeploymentType); + + return initiateDeployment( + context, + transition.nextBlueGreenDeploymentType, + transition.nextState, + null, // Use initialSavepointPath from spec + false); + } + private DeploymentTransition calculateTransition(BlueGreenDeploymentType currentType) { if (BlueGreenDeploymentType.BLUE == currentType) { return new DeploymentTransition( @@ -411,7 +461,10 @@ private UpdateControl handleSpecChangesDuringTransitio context.getDeploymentByType(oppositeDeploymentType) .getMetadata() .getName()); - return patchFlinkDeployment(context, oppositeDeploymentType); + return patchFlinkDeployment( + context, + oppositeDeploymentType, + diffType != BlueGreenDiffType.SAVEPOINT_REDEPLOY); } } diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java index 04c243d43d..a43f3168ac 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java @@ -59,24 +59,31 @@ public BlueGreenDiffType compare() { FlinkDeploymentSpec leftSpec = left.getTemplate().getSpec(); FlinkDeploymentSpec rightSpec = right.getTemplate().getSpec(); - // Used in Case 2 & 3: Delegate to ReflectiveDiffBuilder for nested spec comparison + // Used in Case 2, 3 & 4: Delegate to ReflectiveDiffBuilder for nested spec comparison // Calculate diffResult before comparison to apply in-place removal of ignored fields DiffResult diffResult = new ReflectiveDiffBuilder<>(deploymentMode, leftSpec, rightSpec).build(); - // Case 1: FlinkDeploymentSpecs are identical + // Extract the diff type from ReflectiveDiffBuilder result + DiffType diffType = diffResult.getType(); + + // Case 1: Check for savepoint redeploy first (savepointRedeployNonce changed) + // This takes precedence as it indicates user wants to redeploy from their + // initialSavepointPath + if (diffType == DiffType.SAVEPOINT_REDEPLOY) { + return BlueGreenDiffType.SAVEPOINT_REDEPLOY; + } + + // Case 2: FlinkDeploymentSpecs are identical if (leftSpec.equals(rightSpec)) { return BlueGreenDiffType.IGNORE; } - // Extract the diff type from ReflectiveDiffBuilder result - DiffType diffType = diffResult.getType(); - - // Case 2: ReflectiveDiffBuilder returns IGNORE + // Case 3: ReflectiveDiffBuilder returns IGNORE if (diffType == DiffType.IGNORE) { return BlueGreenDiffType.PATCH_CHILD; } else { - // Case 3: ReflectiveDiffBuilder returns anything else map it to TRANSITION as well + // Case 4: ReflectiveDiffBuilder returns anything else map it to TRANSITION return BlueGreenDiffType.TRANSITION; } } 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 98344a9ba8..2719590884 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 @@ -332,7 +332,12 @@ public static FlinkDeployment prepareFlinkDeployment( "spec", FlinkBlueGreenDeploymentSpec.class); - // The Blue/Green initialSavepointPath is only used for first-time deployments + // Determine which savepoint/checkpoint to restore from: + // 1. First-time deployments: use initialSavepointPath from spec (if set) + // 2. Normal transitions: use lastCheckpoint from previous deployment + // 3. Redeploy scenarios (lastCheckpoint is null): use initialSavepointPath from spec + // - savepointRedeployNonce changed + // - upgradeMode is STATELESS if (isFirstDeployment) { String initialSavepointPath = spec.getTemplate().getSpec().getJob().getInitialSavepointPath(); @@ -346,6 +351,16 @@ public static FlinkDeployment prepareFlinkDeployment( String location = lastCheckpoint.getLocation().replace("file:", ""); LOG.info("Using Blue/Green savepoint/checkpoint: " + location); spec.getTemplate().getSpec().getJob().setInitialSavepointPath(location); + } else { + String initialSavepointPath = + spec.getTemplate().getSpec().getJob().getInitialSavepointPath(); + if (initialSavepointPath != null && !initialSavepointPath.isEmpty()) { + LOG.info( + "Using user-specified initialSavepointPath for redeploy: {}", + initialSavepointPath); + } else { + LOG.info("Starting fresh with no savepoint restoration"); + } } flinkDeployment.setSpec(spec.getTemplate().getSpec()); 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 229406bbc9..8240b8c1b3 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 @@ -189,6 +189,68 @@ private TestingFlinkBlueGreenDeploymentController.BlueGreenReconciliationResult return rs; } + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifySavepointRedeployNonceTriggersTransitionWithInitialSavepointPath( + FlinkVersion flinkVersion) throws Exception { + // Start with SAVEPOINT upgrade mode so normally a savepoint would be taken + var blueGreenDeployment = + buildSessionCluster( + TEST_DEPLOYMENT_NAME, + TEST_NAMESPACE, + flinkVersion, + null, + UpgradeMode.SAVEPOINT); + var rs = executeBasicDeployment(flinkVersion, blueGreenDeployment, false, null); + + // Set initialSavepointPath and bump savepointRedeployNonce + String userSpecifiedSavepoint = "s3://bucket/my-specific-savepoint"; + rs.deployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setInitialSavepointPath(userSpecifiedSavepoint); + rs.deployment.getSpec().getTemplate().getSpec().getJob().setSavepointRedeployNonce(12345L); + rs.deployment + .getSpec() + .getConfiguration() + .put(DEPLOYMENT_DELETION_DELAY.key(), String.valueOf(ALT_DELETION_DELAY_VALUE)); + kubernetesClient.resource(rs.deployment).createOrReplace(); + + // Reconcile - should skip savepointing and go directly to transition + rs = reconcile(rs.deployment); + + // Verify: Should be TRANSITIONING_TO_GREEN (not SAVEPOINTING_BLUE) + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_GREEN, + rs.reconciledStatus.getBlueGreenState()); + + // Verify: Green deployment should use the user-specified initialSavepointPath + var flinkDeployments = getFlinkDeployments(); + assertEquals(2, flinkDeployments.size()); + assertEquals( + userSpecifiedSavepoint, + flinkDeployments.get(1).getSpec().getJob().getInitialSavepointPath()); + + // Complete the transition + simulateSuccessfulJobStart(flinkDeployments.get(1)); + rs = reconcile(rs.deployment); + + // Wait for deletion delay + Thread.sleep(rs.updateControl.getScheduleDelay().get()); + reconcile(rs.deployment); + + // Verify Green is now active + flinkDeployments = getFlinkDeployments(); + assertEquals(1, flinkDeployments.size()); + + rs = reconcile(rs.deployment); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState()); + } + @ParameterizedTest @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") public void verifyFailureBeforeTransition(FlinkVersion flinkVersion) throws Exception { diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java index 9fe6b8bd9b..272730a088 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java @@ -201,13 +201,28 @@ public void testTransitionForNestedSpecDifference() { FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec(); FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); - // Change nested spec property - setSavepointRedeployNonce triggers TRANSITION + // Change nested spec property - setSavepointRedeployNonce now triggers SAVEPOINT_REDEPLOY spec2.getTemplate().getSpec().getJob().setSavepointRedeployNonce(12345L); FlinkBlueGreenDeploymentSpecDiff diff = new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); - assertEquals(BlueGreenDiffType.TRANSITION, diff.compare()); + assertEquals(BlueGreenDiffType.SAVEPOINT_REDEPLOY, diff.compare()); + } + + @Test + public void testSavepointRedeployForNonceChangeWithJarUpdate() { + FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec(); + FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); + + // Nonce change with additional jarURI change - SAVEPOINT_REDEPLOY should still be detected + spec2.getTemplate().getSpec().getJob().setJarURI("local:///opt/flink/examples/other.jar"); + spec2.getTemplate().getSpec().getJob().setSavepointRedeployNonce(12345L); + + FlinkBlueGreenDeploymentSpecDiff diff = + new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); + + assertEquals(BlueGreenDiffType.SAVEPOINT_REDEPLOY, diff.compare()); } @Test @@ -281,8 +296,8 @@ public void testTransitionForTopLevelAndNestedDifferences() { FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); // Change both top-level (configuration) and nested spec - // With new logic, only nested spec changes matter - setSavepointRedeployNonce triggers - // TRANSITION + // With new logic, only nested spec changes matter - setSavepointRedeployNonce now + // triggers SAVEPOINT_REDEPLOY Map config = new HashMap<>(); config.put("custom.config", "different-value"); spec2.setConfiguration(config); @@ -291,7 +306,39 @@ public void testTransitionForTopLevelAndNestedDifferences() { FlinkBlueGreenDeploymentSpecDiff diff = new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); - assertEquals(BlueGreenDiffType.TRANSITION, diff.compare()); + assertEquals(BlueGreenDiffType.SAVEPOINT_REDEPLOY, diff.compare()); + } + + @Test + public void testSavepointRedeployTakesPrecedenceOverScale() { + FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec(); + FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); + + // Change both nonce (SAVEPOINT_REDEPLOY) AND parallelism (SCALE) + // SAVEPOINT_REDEPLOY should take precedence + spec2.getTemplate().getSpec().getJob().setSavepointRedeployNonce(12345L); + spec2.getTemplate().getSpec().getJob().setParallelism(10); + + FlinkBlueGreenDeploymentSpecDiff diff = + new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); + + assertEquals(BlueGreenDiffType.SAVEPOINT_REDEPLOY, diff.compare()); + } + + @Test + public void testSavepointRedeployTakesPrecedenceOverUpgrade() { + FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec(); + FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); + + // Change both nonce (SAVEPOINT_REDEPLOY) AND Flink version (UPGRADE) + // SAVEPOINT_REDEPLOY should take precedence + spec2.getTemplate().getSpec().getJob().setSavepointRedeployNonce(12345L); + spec2.getTemplate().getSpec().setFlinkVersion(FlinkVersion.v1_17); + + FlinkBlueGreenDeploymentSpecDiff diff = + new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); + + assertEquals(BlueGreenDiffType.SAVEPOINT_REDEPLOY, diff.compare()); } @Test diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtilsTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtilsTest.java index 859d19f5d3..a898662ec8 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtilsTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/bluegreen/BlueGreenUtilsTest.java @@ -27,6 +27,9 @@ import org.apache.flink.kubernetes.operator.api.spec.JobSpec; import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode; import org.apache.flink.kubernetes.operator.api.status.FlinkBlueGreenDeploymentStatus; +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.controller.bluegreen.BlueGreenContext; @@ -38,6 +41,9 @@ import java.util.UUID; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; /** Tests for {@link BlueGreenUtils}. */ public class BlueGreenUtilsTest { @@ -83,6 +89,141 @@ public void testPrepareFlinkDeploymentWithoutNameReplacement() { assertEquals(parentDeploymentName + ".jm", resultFlinkConfig.get("metrics.scope.jm")); } + @Test + public void testSavepointRequiredBasedOnUpgradeMode() { + // SAVEPOINT mode requires savepoint + FlinkBlueGreenDeployment bgDeployment = + buildBlueGreenDeployment("test-app", TEST_NAMESPACE); + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setUpgradeMode(UpgradeMode.SAVEPOINT); + BlueGreenContext context = createContext(bgDeployment); + assertTrue(BlueGreenUtils.isSavepointRequired(context)); + + // LAST_STATE mode requires savepoint + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setUpgradeMode(UpgradeMode.LAST_STATE); + assertTrue(BlueGreenUtils.isSavepointRequired(context)); + + // STATELESS mode does not require savepoint + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setUpgradeMode(UpgradeMode.STATELESS); + assertFalse(BlueGreenUtils.isSavepointRequired(context)); + } + + @Test + public void testPrepareFlinkDeploymentStatelessInitialSavepointPath() { + // Setup: STATELESS mode with initialSavepointPath set + FlinkBlueGreenDeployment bgDeployment = + buildBlueGreenDeployment("test-app", TEST_NAMESPACE); + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setUpgradeMode(UpgradeMode.STATELESS); + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setInitialSavepointPath("s3://bucket/savepoint-xyz"); + + BlueGreenContext context = createContext(bgDeployment); + + // Act: Prepare deployment with null lastCheckpoint (STATELESS transition) + FlinkDeployment result = + BlueGreenUtils.prepareFlinkDeployment( + context, + BlueGreenDeploymentType.GREEN, + null, // No lastCheckpoint + false, // Not first deployment + bgDeployment.getMetadata()); + + // Assert: initialSavepointPath should be used for STATELESS + assertNotNull(result.getSpec().getJob().getInitialSavepointPath()); + } + + @Test + public void testNullLastCheckpointUsesInitialSavepointPath() { + // lastCheckpoint=null -> use initialSavepointPath from spec + FlinkBlueGreenDeployment bgDeployment = + buildBlueGreenDeployment("test-app", TEST_NAMESPACE); + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setUpgradeMode(UpgradeMode.SAVEPOINT); + String initialPath = "s3://bucket/user-specified-savepoint"; + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setInitialSavepointPath(initialPath); + + BlueGreenContext context = createContext(bgDeployment); + + FlinkDeployment result = + BlueGreenUtils.prepareFlinkDeployment( + context, + BlueGreenDeploymentType.GREEN, + null, // null = nonce changed, no new savepoint taken + false, + bgDeployment.getMetadata()); + + assertEquals(initialPath, result.getSpec().getJob().getInitialSavepointPath()); + } + + @Test + public void testNormalTransitionUsesFreshSavepoint() { + // Normal transition → take fresh savepoint from running job → use that, not + // initialSavepointPath + FlinkBlueGreenDeployment bgDeployment = + buildBlueGreenDeployment("test-app", TEST_NAMESPACE); + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setUpgradeMode(UpgradeMode.SAVEPOINT); + bgDeployment + .getSpec() + .getTemplate() + .getSpec() + .getJob() + .setInitialSavepointPath("s3://bucket/ignored"); + + BlueGreenContext context = createContext(bgDeployment); + + String freshSavepoint = "s3://bucket/fresh-savepoint-from-running-job"; + Savepoint triggered = + Savepoint.of( + freshSavepoint, SnapshotTriggerType.UPGRADE, SavepointFormatType.CANONICAL); + + FlinkDeployment result = + BlueGreenUtils.prepareFlinkDeployment( + context, + BlueGreenDeploymentType.GREEN, + triggered, // Fresh savepoint provided + false, + bgDeployment.getMetadata()); + + assertEquals(freshSavepoint, result.getSpec().getJob().getInitialSavepointPath()); + } + private static FlinkBlueGreenDeployment buildBlueGreenDeployment( String name, String namespace) { var deployment = new FlinkBlueGreenDeployment(); From 4847dfd422dd4d8f0ca12520aed04082b7f3e486 Mon Sep 17 00:00:00 2001 From: James Kan Date: Tue, 13 Jan 2026 19:52:19 -0800 Subject: [PATCH 8/8] [FLINK-38915] FlinkBlueGreenDeplomynet in place suspension handler --- .../api/bluegreen/BlueGreenDiffType.java | 14 +- .../bluegreen/BlueGreenDeploymentService.java | 84 +++++- .../InitializingBlueStateHandler.java | 11 + .../FlinkBlueGreenDeploymentSpecDiff.java | 27 ++ ...linkBlueGreenDeploymentControllerTest.java | 242 ++++++++++++++++++ .../FlinkBlueGreenDeploymentSpecDiffTest.java | 24 ++ 6 files changed, 398 insertions(+), 4 deletions(-) diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java index e456a89133..97a83ec212 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/api/bluegreen/BlueGreenDiffType.java @@ -32,5 +32,17 @@ public enum BlueGreenDiffType { * Full redeploy from user-specified savepoint. Triggered when savepointRedeployNonce changes. * Uses the initialSavepointPath from spec instead of taking a new savepoint. */ - SAVEPOINT_REDEPLOY + SAVEPOINT_REDEPLOY, + + /** + * In-place suspension. Triggered when job.state changes from RUNNING to SUSPENDED. Suspends the + * currently active child without creating a new deployment. + */ + SUSPEND, + + /** + * Resume from suspension. Triggered when job.state changes from SUSPENDED to RUNNING. Spins up + * the child with the current (potentially updated) spec. + */ + RESUME } 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 132b05e27f..33cd868b0b 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,6 +22,7 @@ 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.lifecycle.ResourceLifecycleState; 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; @@ -110,13 +111,38 @@ public UpdateControl initiateDeployment( */ public UpdateControl checkAndInitiateDeployment( BlueGreenContext context, BlueGreenDeploymentType currentBlueGreenDeploymentType) { + BlueGreenDiffType specDiff = getSpecDiff(context); if (specDiff != BlueGreenDiffType.IGNORE) { FlinkDeployment currentFlinkDeployment = context.getDeploymentByType(currentBlueGreenDeploymentType); - if (isFlinkDeploymentReady(currentFlinkDeployment)) { + if (specDiff == BlueGreenDiffType.SUSPEND && currentFlinkDeployment != null) { + setLastReconciledSpec(context); + LOG.info( + "In-place suspension for '{}'", + currentFlinkDeployment.getMetadata().getName()); + return patchFlinkDeployment(context, currentBlueGreenDeploymentType); + } + + if (specDiff == BlueGreenDiffType.RESUME && currentFlinkDeployment != null) { + setLastReconciledSpec(context); + LOG.info( + "In-place resume for '{}'", currentFlinkDeployment.getMetadata().getName()); + return patchFlinkDeployment(context, currentBlueGreenDeploymentType); + } + + // Check if child is currently suspended - if so, just patch specs without restart + if (isChildSuspended(currentFlinkDeployment)) { + setLastReconciledSpec(context); + LOG.info( + "Spec change while suspended for '{}'", + currentFlinkDeployment.getMetadata().getName()); + return patchFlinkDeployment(context, currentBlueGreenDeploymentType); + } + + if (currentFlinkDeployment != null && isFlinkDeploymentReady(currentFlinkDeployment)) { if (specDiff == BlueGreenDiffType.TRANSITION) { boolean savepointTriggered = false; try { @@ -173,12 +199,16 @@ public UpdateControl checkAndInitiateDeployment( } else { if (context.getDeploymentStatus().getJobStatus().getState() != JobStatus.FAILING) { setLastReconciledSpec(context); + var childName = + currentFlinkDeployment != null + ? currentFlinkDeployment.getMetadata().getName() + : "null"; var error = String.format( "Transition to %s not possible, current Flink Deployment '%s' is not READY. FAILING '%s'", calculateTransition(currentBlueGreenDeploymentType) .nextBlueGreenDeploymentType, - currentFlinkDeployment.getMetadata().getName(), + childName, context.getBgDeployment().getMetadata().getName()); return markDeploymentFailing(context, error); } @@ -188,6 +218,16 @@ public UpdateControl checkAndInitiateDeployment( return UpdateControl.noUpdate(); } + private boolean isChildSuspended(FlinkDeployment deployment) { + if (deployment == null || deployment.getSpec() == null) { + return false; + } + var job = deployment.getSpec().getJob(); + return job != null + && job.getState() + == org.apache.flink.kubernetes.operator.api.spec.JobState.SUSPENDED; + } + private UpdateControl patchFlinkDeployment( BlueGreenContext context, BlueGreenDeploymentType blueGreenDeploymentTypeToPatch) { return patchFlinkDeployment(context, blueGreenDeploymentTypeToPatch, true); @@ -435,6 +475,16 @@ public UpdateControl monitorTransition( TransitionState transitionState = determineTransitionState(context, currentBlueGreenDeploymentType); + if (isChildSuspended(transitionState.nextDeployment)) { + if (transitionState.nextDeployment.getStatus().getLifecycleState() + == ResourceLifecycleState.SUSPENDED) { + return finalizeSuspendedDeployment(context, transitionState.nextState); + } else { + return shouldWeAbort( + context, transitionState.nextDeployment, transitionState.nextState); + } + } + if (isFlinkDeploymentReady(transitionState.nextDeployment)) { return shouldWeDelete( context, @@ -447,11 +497,36 @@ public UpdateControl monitorTransition( } } + private UpdateControl finalizeSuspendedDeployment( + BlueGreenContext context, FlinkBlueGreenDeploymentState nextState) { + + LOG.info( + "Finalizing suspended deployment '{}' to {} state", + context.getDeploymentName(), + nextState); + + context.getDeploymentStatus().setDeploymentReadyTimestamp(millisToInstantStr(0)); + context.getDeploymentStatus().setAbortTimestamp(millisToInstantStr(0)); + context.getDeploymentStatus().setSavepointTriggerId(null); + + return patchStatusUpdateControl(context, nextState, JobStatus.SUSPENDED, null) + .rescheduleAfter(0); + } + 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; + } + if (diffType != BlueGreenDiffType.IGNORE) { setLastReconciledSpec(context); var oppositeDeploymentType = @@ -658,7 +733,10 @@ public UpdateControl finalizeBlueGreenDeployment( context.getDeploymentStatus().setAbortTimestamp(millisToInstantStr(0)); context.getDeploymentStatus().setSavepointTriggerId(null); - return patchStatusUpdateControl(context, nextState, JobStatus.RUNNING, null); + // Finalize status and reschedule immediately so any pending spec changes + // (e.g., suspend requested during transition) are picked up on next reconcile + return patchStatusUpdateControl(context, nextState, JobStatus.RUNNING, null) + .rescheduleAfter(0); } // ==================== Common Utility Methods ==================== diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/handlers/InitializingBlueStateHandler.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/handlers/InitializingBlueStateHandler.java index f2882d46f7..8319314fb6 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/handlers/InitializingBlueStateHandler.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/controller/bluegreen/handlers/InitializingBlueStateHandler.java @@ -20,6 +20,7 @@ import org.apache.flink.api.common.JobStatus; import org.apache.flink.kubernetes.operator.api.FlinkBlueGreenDeployment; import org.apache.flink.kubernetes.operator.api.bluegreen.BlueGreenDeploymentType; +import org.apache.flink.kubernetes.operator.api.spec.JobState; import org.apache.flink.kubernetes.operator.api.status.FlinkBlueGreenDeploymentState; import org.apache.flink.kubernetes.operator.api.status.FlinkBlueGreenDeploymentStatus; import org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenContext; @@ -27,6 +28,7 @@ import io.javaoperatorsdk.operator.api.reconciler.UpdateControl; +import static org.apache.flink.kubernetes.operator.controller.bluegreen.BlueGreenDeploymentService.patchStatusUpdateControl; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.hasSpecChanged; import static org.apache.flink.kubernetes.operator.utils.bluegreen.BlueGreenUtils.setLastReconciledSpec; @@ -41,6 +43,15 @@ public InitializingBlueStateHandler(BlueGreenDeploymentService deploymentService public UpdateControl handle(BlueGreenContext context) { FlinkBlueGreenDeploymentStatus deploymentStatus = context.getDeploymentStatus(); + // Block initial deployment if job.state is SUSPENDED - user must start with RUNNING + var jobSpec = context.getBgDeployment().getSpec().getTemplate().getSpec().getJob(); + if (jobSpec != null && jobSpec.getState() == JobState.SUSPENDED) { + LOG.info( + "Blocking initial deployment '{}' - job.state is SUSPENDED, waiting for RUNNING", + context.getBgDeployment().getMetadata().getName()); + return patchStatusUpdateControl(context, null, JobStatus.SUSPENDED, null); + } + // Deploy only if this is the initial deployment (no previous spec exists) // or if we're recovering from a failure and the spec has changed since the last attempt if (deploymentStatus.getLastReconciledSpec() == null diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java index a43f3168ac..3e5dba77cb 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiff.java @@ -22,6 +22,7 @@ import org.apache.flink.kubernetes.operator.api.diff.DiffType; 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 lombok.NonNull; @@ -59,6 +60,19 @@ public BlueGreenDiffType compare() { FlinkDeploymentSpec leftSpec = left.getTemplate().getSpec(); FlinkDeploymentSpec rightSpec = right.getTemplate().getSpec(); + // Check for suspend/resume state changes first - these take highest precedence + // for in-place suspension handling + JobState leftJobState = getJobState(leftSpec); + JobState rightJobState = getJobState(rightSpec); + + if (leftJobState != rightJobState) { + if (rightJobState == JobState.SUSPENDED) { + return BlueGreenDiffType.SUSPEND; + } else if (leftJobState == JobState.SUSPENDED && rightJobState == JobState.RUNNING) { + return BlueGreenDiffType.RESUME; + } + } + // Used in Case 2, 3 & 4: Delegate to ReflectiveDiffBuilder for nested spec comparison // Calculate diffResult before comparison to apply in-place removal of ignored fields DiffResult diffResult = @@ -88,6 +102,19 @@ public BlueGreenDiffType compare() { } } + /** + * Gets the job state from the spec, defaulting to RUNNING if not set. + * + * @param spec the FlinkDeploymentSpec + * @return the job state, or RUNNING if job or state is null + */ + private JobState getJobState(FlinkDeploymentSpec spec) { + if (spec.getJob() == null || spec.getJob().getState() == null) { + return JobState.RUNNING; + } + return spec.getJob().getState(); + } + /** * Validates that the specs and their nested components are not null. Throws * IllegalArgumentException if any required component is 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 8240b8c1b3..cf99b2e99e 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 @@ -251,6 +251,168 @@ public void verifySavepointRedeployNonceTriggersTransitionWithInitialSavepointPa rs.reconciledStatus.getBlueGreenState()); } + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifySuspendAndResumeInPlace(FlinkVersion flinkVersion) throws Exception { + var rs = setupActiveBlueDeployment(flinkVersion); + assertEquals(1, getFlinkDeployments().size()); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should start as ACTIVE_BLUE"); + + // Suspend in-place with spec change (parallelism bump) + var deployment = rs.deployment; + deployment.getSpec().getTemplate().getSpec().getJob().setState(JobState.SUSPENDED); + deployment.getSpec().getTemplate().getSpec().getJob().setParallelism(5); + kubernetesClient.resource(deployment).createOrReplace(); + + rs = reconcile(deployment); + deployment = rs.deployment; + // Verify BG goes through TRANSITIONING_TO_BLUE state (unified patch mechanism) + assertEquals( + JobStatus.RECONCILING, + rs.reconciledStatus.getJobStatus().getState(), + "BG status should be RECONCILING after suspend initiated"); + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should go through TRANSITIONING_TO_BLUE during suspend"); + + var flinkDeployments = getFlinkDeployments(); + assertEquals(1, flinkDeployments.size(), "Suspend should keep single child"); + var child = flinkDeployments.get(0); + assertTrue( + child.getMetadata().getName().endsWith("-blue"), "Child should be blue deployment"); + assertEquals( + JobState.SUSPENDED, + child.getSpec().getJob().getState(), + "Child should be suspended"); + assertEquals(5, child.getSpec().getJob().getParallelism(), "Spec change should be applied"); + + // Simulate child becoming suspended (lifecycleState is computed from lastReconciledSpec) + simulateSuccessfulSuspend(child); + rs = reconcile(deployment); + deployment = rs.deployment; + assertEquals(JobStatus.SUSPENDED, rs.reconciledStatus.getJobStatus().getState()); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should finalize to ACTIVE_BLUE when suspended"); + + // Resume in-place with another spec change + deployment.getSpec().getTemplate().getSpec().getJob().setState(JobState.RUNNING); + deployment.getSpec().getTemplate().getSpec().getJob().setParallelism(6); + kubernetesClient.resource(deployment).createOrReplace(); + + rs = reconcile(deployment); + deployment = rs.deployment; + // Resume also goes through TRANSITIONING_TO_BLUE (unified patch mechanism) + assertEquals( + FlinkBlueGreenDeploymentState.TRANSITIONING_TO_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should go through TRANSITIONING_TO_BLUE during resume"); + // Persist the BG status update to kubernetes for the next reconcile + kubernetesClient.resource(deployment).updateStatus(); + + flinkDeployments = getFlinkDeployments(); + assertEquals(1, flinkDeployments.size(), "Resume should keep single child"); + child = flinkDeployments.get(0); + assertTrue( + child.getMetadata().getName().endsWith("-blue"), + "Child should still be blue deployment"); + assertEquals(JobState.RUNNING, child.getSpec().getJob().getState(), "Child should resume"); + assertEquals(6, child.getSpec().getJob().getParallelism(), "Latest spec should be applied"); + + // Simulate child becoming running and verify BG status syncs + simulateSuccessfulJobStart(child); + rs = reconcile(deployment); + assertEquals(JobStatus.RUNNING, rs.reconciledStatus.getJobStatus().getState()); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_BLUE, + rs.reconciledStatus.getBlueGreenState(), + "Should finalize to ACTIVE_BLUE after resume"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifySuspendDuringTransitionIsBlockedUntilComplete(FlinkVersion flinkVersion) + throws Exception { + // Start with a running Blue deployment with longer abort grace period + var blueGreenDeployment = + buildSessionCluster( + TEST_DEPLOYMENT_NAME, + TEST_NAMESPACE, + flinkVersion, + null, + UpgradeMode.STATELESS); // STATELESS to skip savepointing + // Set longer abort grace period BEFORE deployment to avoid abort during transition test + blueGreenDeployment.getSpec().getConfiguration().put(ABORT_GRACE_PERIOD.key(), "60000"); + // Use zero deletion delay to avoid timing-based flakiness during transition completion + blueGreenDeployment.getSpec().getConfiguration().put(DEPLOYMENT_DELETION_DELAY.key(), "0"); + var rs = executeBasicDeployment(flinkVersion, blueGreenDeployment, false, null); + + // === TRIGGER TRANSITION === + rs.deployment.getSpec().getTemplate().getSpec().getJob().setParallelism(5); + kubernetesClient.resource(rs.deployment).createOrReplace(); + + // Reconcile - start transition to GREEN + rs = reconcile(rs.deployment); + + // GREEN deployment should be created + var flinkDeployments = getFlinkDeployments(); + assertEquals(2, flinkDeployments.size()); + + // === SUSPEND DURING TRANSITION === + rs.deployment.getSpec().getTemplate().getSpec().getJob().setState(JobState.SUSPENDED); + kubernetesClient.resource(rs.deployment).createOrReplace(); + + // Reconcile - suspend should be BLOCKED, transition continues + rs = reconcile(rs.deployment); + + // GREEN deployment should NOT be suspended yet + flinkDeployments = getFlinkDeployments(); + var greenDeployment = + flinkDeployments.stream() + .filter(fd -> fd.getMetadata().getName().endsWith("-green")) + .findFirst() + .orElseThrow(); + assertEquals(JobState.RUNNING, greenDeployment.getSpec().getJob().getState()); + + // === TRANSITION COMPLETES === + // Simulate GREEN becoming ready + simulateSuccessfulJobStart(greenDeployment); + + // Reconcile to complete transition and process pending suspend + // (sets timestamp, deletes BLUE, finalizes to ACTIVE_GREEN, then processes suspend + // via TRANSITIONING_TO_GREEN) + for (int i = 0; i < 5; i++) { + rs = reconcile(rs.deployment); + } + // Persist the BG status update to kubernetes for the next reconcile + kubernetesClient.resource(rs.deployment).updateStatus(); + + // GREEN should now be suspended in spec (via patchFlinkDeployment) + flinkDeployments = getFlinkDeployments(); + greenDeployment = + flinkDeployments.stream() + .filter(fd -> fd.getMetadata().getName().endsWith("-green")) + .findFirst() + .orElseThrow(); + assertEquals(JobState.SUSPENDED, greenDeployment.getSpec().getJob().getState()); + + // Simulate GREEN child becoming suspended (lifecycleState is computed from + // lastReconciledSpec) + simulateSuccessfulSuspend(greenDeployment); + + // Reconcile - finalizeSuspendedDeployment should set BG job status to SUSPENDED + rs = reconcile(rs.deployment); + assertEquals(JobStatus.SUSPENDED, rs.reconciledStatus.getJobStatus().getState()); + assertEquals( + FlinkBlueGreenDeploymentState.ACTIVE_GREEN, + rs.reconciledStatus.getBlueGreenState()); + } + @ParameterizedTest @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") public void verifyFailureBeforeTransition(FlinkVersion flinkVersion) throws Exception { @@ -654,6 +816,76 @@ public void verifySavepointFetchFailureWithDifferentErrors( assertEquals(1, flinkDeployments.size()); } + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifySuspendWhenChildNotReadyPatchesWithoutFailing(FlinkVersion flinkVersion) + throws Exception { + var rs = setupActiveBlueDeployment(flinkVersion); + + // Mark the active child as not ready + var child = getFlinkDeployments().get(0); + child.getStatus().getJobStatus().setState(JobStatus.RECONCILING); + kubernetesClient.resource(child).update(); + + // Suspend and bump parallelism + var deployment = rs.deployment; + deployment.getSpec().getTemplate().getSpec().getJob().setState(JobState.SUSPENDED); + deployment.getSpec().getTemplate().getSpec().getJob().setParallelism(9); + kubernetesClient.resource(deployment).createOrReplace(); + + rs = reconcile(deployment); + + var flinkDeployments = getFlinkDeployments(); + assertEquals(1, flinkDeployments.size(), "Should not create new child"); + child = flinkDeployments.get(0); + assertEquals(JobState.SUSPENDED, child.getSpec().getJob().getState()); + assertEquals(9, child.getSpec().getJob().getParallelism(), "Spec change should apply"); + assertNotEquals( + JobStatus.FAILING, + rs.reconciledStatus.getJobStatus().getState(), + "Should not mark parent failing when suspending not-ready child"); + } + + @ParameterizedTest + @MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions") + public void verifyInitialSuspendedIsBlockedThenDeploysOnRunning(FlinkVersion flinkVersion) + throws Exception { + var bg = + buildSessionCluster( + TEST_DEPLOYMENT_NAME, + TEST_NAMESPACE, + flinkVersion, + null, + UpgradeMode.STATELESS); + bg.getSpec().getTemplate().getSpec().getJob().setState(JobState.SUSPENDED); + kubernetesClient.resource(bg).createOrReplace(); + + // First reconcile initializes status to INITIALIZING_BLUE + var rs = reconcile(bg); + // Second reconcile: handler detects SUSPENDED and blocks + rs = reconcile(rs.deployment); + + // Block child creation when initial state is SUSPENDED + assertEquals(0, getFlinkDeployments().size(), "No child should be created"); + assertEquals( + JobStatus.SUSPENDED, + rs.reconciledStatus.getJobStatus().getState(), + "Job status should be SUSPENDED when initial deployment is blocked"); + + // Flip to RUNNING and reconcile again + bg = rs.deployment; + bg.getSpec().getTemplate().getSpec().getJob().setState(JobState.RUNNING); + kubernetesClient.resource(bg).createOrReplace(); + + rs = reconcile(bg); + + assertEquals(1, getFlinkDeployments().size(), "Child should be created after fix"); + assertEquals( + JobStatus.RECONCILING, + rs.reconciledStatus.getJobStatus().getState(), + "Job status should be RECONCILING after deployment initiated"); + } + // ==================== Parameterized Test Inputs ==================== static Stream savepointErrorProvider() { @@ -1196,6 +1428,16 @@ private void simulateSuccessfulJobStart(FlinkDeployment deployment) { kubernetesClient.resource(deployment).update(); } + private void simulateSuccessfulSuspend(FlinkDeployment deployment) { + deployment.getStatus().getJobStatus().setState(JobStatus.FINISHED); + deployment.getStatus().getReconciliationStatus().setState(ReconciliationState.DEPLOYED); + deployment + .getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(deployment.getSpec(), deployment); + kubernetesClient.resource(deployment).update(); + } + private void simulateJobFailure(FlinkDeployment deployment) { deployment.getStatus().getJobStatus().setState(JobStatus.RECONCILING); deployment.getStatus().getReconciliationStatus().setState(ReconciliationState.UPGRADING); diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java index 272730a088..1e60b83f4b 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/diff/FlinkBlueGreenDeploymentSpecDiffTest.java @@ -134,6 +134,30 @@ public void testIgnoreForConfigurationDifference() { assertEquals(BlueGreenDiffType.IGNORE, diff.compare()); } + @Test + public void testSuspendOnJobStateChange() { + FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec(); // RUNNING default + FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); + spec2.getTemplate().getSpec().getJob().setState(JobState.SUSPENDED); + + FlinkBlueGreenDeploymentSpecDiff diff = + new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); + + assertEquals(BlueGreenDiffType.SUSPEND, diff.compare()); + } + + @Test + public void testResumeOnJobStateChange() { + FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec(); + spec1.getTemplate().getSpec().getJob().setState(JobState.SUSPENDED); + FlinkBlueGreenDeploymentSpec spec2 = createBasicSpec(); // RUNNING default + + FlinkBlueGreenDeploymentSpecDiff diff = + new FlinkBlueGreenDeploymentSpecDiff(DEPLOYMENT_MODE, spec1, spec2); + + assertEquals(BlueGreenDiffType.RESUME, diff.compare()); + } + @Test public void testIgnoreForRootPodTemplateAdditionalProps() { FlinkBlueGreenDeploymentSpec spec1 = createBasicSpec();