context) { } /** - * Resolves the replica count to set in the desired workload spec. + * Resolves the replica count to set in the desired workload spec, and logs it when it differs + * from what the workload has now. *
- * Always returns an explicit value — never null. Returning null would cause - * JOSDK/SSA to omit spec.replicas, and Kubernetes would default it to 1. + * Returns a primitive so the value can never be null: a null spec.replicas would make JOSDK/SSA + * omit the field, and Kubernetes would default it to 1. *
* When autoscaling is enabled: * - On CREATE: returns initialReplicas (minReplicas for the component) @@ -173,7 +175,20 @@ protected R handleCreate(R desired, P primary, Context
context) { *
* When autoscaling is disabled: returns staticReplicas (the spec value). */ - protected Integer resolveReplicaCount(P primary, Context
context, + protected int resolveReplicaCount(P primary, Context
context,
+ AutoscalingSpec autoscaling, int staticReplicas, int initialReplicas) {
+ Optional context,
if (autoscaling == null || !autoscaling.isEnabled()) {
return staticReplicas;
}
- Optional context,
if (managed != null) {
return managed;
}
- // Fallback: operator restarted and MANAGED_REPLICAS is empty — read current value
- R resource = existing.get();
- if (resource instanceof io.fabric8.kubernetes.api.model.apps.Deployment d) {
- return d.getSpec() != null && d.getSpec().getReplicas() != null
- ? d.getSpec().getReplicas() : initialReplicas;
- }
- if (resource instanceof io.fabric8.kubernetes.api.model.apps.StatefulSet s) {
- return s.getSpec() != null && s.getSpec().getReplicas() != null
- ? s.getSpec().getReplicas() : initialReplicas;
- }
- return initialReplicas;
+ // Fallback: operator restarted and MANAGED_REPLICAS is empty — read current value. The
+ // workload exists, so spec.replicas is set unless something wrote it away; initialReplicas
+ // is the floor either way.
+ return Workloads.replicas(existing.get()).orElse(initialReplicas);
}
// First creation: start at minReplicas.
return initialReplicas;
}
-
/**
* Returns the component name for this dependent (used for autoscaler replica lookup).
* Subclasses should override if they manage a workload with autoscaling.
diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveServer2DeploymentDependent.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveServer2DeploymentDependent.java
index 6bc6291fd1fc..a8e8d953dd7a 100644
--- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveServer2DeploymentDependent.java
+++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveServer2DeploymentDependent.java
@@ -232,7 +232,7 @@ protected Deployment desired(HiveCluster hiveCluster,
AutoscalingSpec hs2Autoscaling = hs2.autoscaling();
int initialReplicas = hs2Autoscaling != null && hs2Autoscaling.isEnabled()
? Math.max(1, hs2Autoscaling.minReplicas()) : hs2.replicas();
- Integer replicas = resolveReplicaCount(
+ int replicas = resolveReplicaCount(
hiveCluster, context, hs2Autoscaling, hs2.replicas(), initialReplicas);
Deployment deployment = new DeploymentBuilder()
diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/LlapResourceBuilder.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/LlapResourceBuilder.java
index a2c10d3688c4..71acf0db4814 100644
--- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/LlapResourceBuilder.java
+++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/LlapResourceBuilder.java
@@ -60,6 +60,7 @@
import org.apache.hive.kubernetes.operator.util.HadoopXmlBuilder;
import org.apache.hive.kubernetes.operator.util.HiveConfigBuilder;
import org.apache.hive.kubernetes.operator.util.Labels;
+import org.apache.hive.kubernetes.operator.util.Workloads;
import static org.apache.hive.kubernetes.operator.autoscaling.MetricsScraper.isPodReady;
@@ -74,7 +75,6 @@ public class LlapResourceBuilder
extends HiveDependentResource
+ * Shared by every scale path — the dependents' SSA and the imperative LLAP/TezAM SSAs — so the
+ * operator log reads the same whichever one ran. The caller passes its own logger to keep the
+ * log category pointing at the code that is actually scaling.
+ */
+ public static void logReplicaChange(Logger log, String component, String namespace, String name,
+ HasMetadata current, int desired) {
+ OptionalInt actual = replicas(current);
+ if (actual.isPresent() && actual.getAsInt() == desired) {
+ return;
+ }
+ log.info("Setting replica count for {} {}/{}: {} -> {}", component, namespace, name,
+ actual.isPresent() ? String.valueOf(actual.getAsInt()) : "none", desired);
+ }
+
+ /**
+ * Maps an autoscaler component key to the K8s workload name it drives. Per-LLAP components
+ * carry the LLAP cluster name in their key ("llap-{name}", "tezam-{name}"); everything else
+ * (HS2, Metastore) is a plain "{cluster}-{component}". Kept here so every scale path — the
+ * autoscaler's `patchReplicas`, the imperative LLAP/TezAM SSAs, and the idle-check reads —
+ * resolves the name the same way.
+ *
+ *
+ */
+ public static String nameFor(HiveCluster hc, String component) {
+ String cluster = hc.getMetadata().getName();
+ if (component.startsWith(ConfigUtils.COMPONENT_LLAP + "-")) {
+ String llapName = component.substring(ConfigUtils.COMPONENT_LLAP.length() + 1);
+ return cluster + "-" + llapName;
+ }
+ if (component.startsWith(ConfigUtils.COMPONENT_TEZAM + "-")) {
+ String llapName = component.substring(ConfigUtils.COMPONENT_TEZAM.length() + 1);
+ return cluster + "-" + ConfigUtils.COMPONENT_TEZAM + "-" + llapName;
+ }
+ return cluster + "-" + component;
+ }
+}
diff --git a/packaging/src/kubernetes/src/test/java/org/apache/hive/kubernetes/operator/util/TestWorkloads.java b/packaging/src/kubernetes/src/test/java/org/apache/hive/kubernetes/operator/util/TestWorkloads.java
new file mode 100644
index 000000000000..3d2e6d537796
--- /dev/null
+++ b/packaging/src/kubernetes/src/test/java/org/apache/hive/kubernetes/operator/util/TestWorkloads.java
@@ -0,0 +1,169 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hive.kubernetes.operator.util;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+
+import java.util.OptionalInt;
+
+import io.fabric8.kubernetes.api.model.ConfigMap;
+import io.fabric8.kubernetes.api.model.ConfigMapBuilder;
+import io.fabric8.kubernetes.api.model.ObjectMetaBuilder;
+import io.fabric8.kubernetes.api.model.apps.Deployment;
+import io.fabric8.kubernetes.api.model.apps.DeploymentBuilder;
+import io.fabric8.kubernetes.api.model.apps.StatefulSet;
+import io.fabric8.kubernetes.api.model.apps.StatefulSetBuilder;
+import org.apache.hive.kubernetes.operator.model.HiveCluster;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+import org.slf4j.Logger;
+
+class TestWorkloads {
+
+ @Test
+ void replicasReturnsValueFromDeployment() {
+ Deployment d = new DeploymentBuilder()
+ .withNewSpec().withReplicas(3).endSpec()
+ .build();
+ assertEquals(OptionalInt.of(3), Workloads.replicas(d));
+ }
+
+ @Test
+ void replicasReturnsValueFromStatefulSet() {
+ StatefulSet s = new StatefulSetBuilder()
+ .withNewSpec().withReplicas(5).endSpec()
+ .build();
+ assertEquals(OptionalInt.of(5), Workloads.replicas(s));
+ }
+
+ @Test
+ void replicasIsEmptyWhenDeploymentSpecMissing() {
+ // fabric8 Deployment with no spec set at all
+ Deployment d = new DeploymentBuilder().build();
+ assertFalse(Workloads.replicas(d).isPresent());
+ }
+
+ @Test
+ void replicasIsEmptyWhenStatefulSetSpecMissing() {
+ StatefulSet s = new StatefulSetBuilder().build();
+ assertFalse(Workloads.replicas(s).isPresent());
+ }
+
+ @Test
+ void replicasIsEmptyWhenDeploymentReplicasUnset() {
+ // spec present, but replicas field not set
+ Deployment d = new DeploymentBuilder().withNewSpec().endSpec().build();
+ assertFalse(Workloads.replicas(d).isPresent());
+ }
+
+ @Test
+ void replicasIsEmptyWhenStatefulSetReplicasUnset() {
+ StatefulSet s = new StatefulSetBuilder().withNewSpec().endSpec().build();
+ assertFalse(Workloads.replicas(s).isPresent());
+ }
+
+ @Test
+ void replicasIsEmptyForNonWorkloadResource() {
+ // Anything that's neither a Deployment nor a StatefulSet is empty,
+ // even if it happens to have a "spec" (e.g., a ConfigMap here has none).
+ ConfigMap cm = new ConfigMapBuilder()
+ .withNewMetadata().withName("cm").endMetadata()
+ .build();
+ assertFalse(Workloads.replicas(cm).isPresent());
+ }
+
+ @Test
+ void replicasIsEmptyForNullResource() {
+ assertFalse(Workloads.replicas(null).isPresent());
+ }
+
+ @Test
+ void logReplicaChangeLogsTheDelta() {
+ Logger log = mock(Logger.class);
+ StatefulSet s = new StatefulSetBuilder().withNewSpec().withReplicas(12).endSpec().build();
+
+ Workloads.logReplicaChange(log, "llap", "ns", "hive-llap0", s, 15);
+
+ verify(log).info(anyString(), eq("llap"), eq("ns"), eq("hive-llap0"), eq("12"), eq(15));
+ }
+
+ /** No workload yet: the same line reports the count the first create will set. */
+ @Test
+ void logReplicaChangeLogsNoneAsTheFromValueOnFirstCreate() {
+ Logger log = mock(Logger.class);
+
+ Workloads.logReplicaChange(log, "hiveserver2", "ns", "hive-hiveserver2", null, 2);
+
+ verify(log).info(anyString(), eq("hiveserver2"), eq("ns"), eq("hive-hiveserver2"),
+ eq("none"), eq(2));
+ }
+
+ /** Silence is the signal that nothing is being scaled, so an unchanged count logs nothing. */
+ @Test
+ void logReplicaChangeIsSilentWhenTheCountAlreadyMatches() {
+ Logger log = mock(Logger.class);
+ Deployment d = new DeploymentBuilder().withNewSpec().withReplicas(2).endSpec().build();
+
+ Workloads.logReplicaChange(log, "metastore", "ns", "hive-metastore", d, 2);
+
+ verifyNoInteractions(log);
+ }
+
+ /**
+ * Covers the three branches of {@link Workloads#nameFor}:
+ *
+ *
+ * The last two rows exercise the "no dash after the prefix" fall-through: a bare
+ * {@code "llap"} or {@code "tezam"} isn't a per-cluster component key and must land in the
+ * plain-{cluster}-{component} branch, not be treated as an empty LLAP name.
+ */
+ @ParameterizedTest
+ @CsvSource({
+ // component, expected workload name
+ "llap-llap0, hive-llap0",
+ "llap-my-llap-cluster, hive-my-llap-cluster",
+ "tezam-llap0, hive-tezam-llap0",
+ "tezam-my-llap-cluster, hive-tezam-my-llap-cluster",
+ "hiveserver2, hive-hiveserver2",
+ "metastore, hive-metastore",
+ "llap, hive-llap",
+ "tezam, hive-tezam",
+ })
+ void nameForMapsComponentKeyToWorkloadName(String component, String expected) {
+ HiveCluster hc = hiveCluster("hive");
+ assertEquals(expected, Workloads.nameFor(hc, component));
+ }
+
+ private static HiveCluster hiveCluster(String name) {
+ HiveCluster hc = new HiveCluster();
+ hc.setMetadata(new ObjectMetaBuilder().withName(name).build());
+ return hc;
+ }
+}