From f403e3c96c989494a0cf2f5f0506060d1db70a98 Mon Sep 17 00:00:00 2001 From: Laszlo Bodor Date: Thu, 10 Sep 2026 11:27:15 +0200 Subject: [PATCH 1/3] HIVE-30033: Hive K8s operator: log replica changes --- .../dependent/HiveDependentResource.java | 47 ++++++++++++++----- .../reconciler/HiveClusterReconciler.java | 29 ++++++++++++ .../kubernetes/operator/util/Workloads.java | 47 +++++++++++++++++++ 3 files changed, 112 insertions(+), 11 deletions(-) create mode 100644 packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java index a64440acab8e..daad1504747a 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java @@ -62,6 +62,7 @@ import org.apache.hive.kubernetes.operator.model.spec.SecretKeyRef; import org.apache.hive.kubernetes.operator.model.spec.ProbeSpec; import org.apache.hive.kubernetes.operator.util.ConfigUtils; +import org.apache.hive.kubernetes.operator.util.Workloads; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -175,6 +176,15 @@ protected R handleCreate(R desired, P primary, Context

context) { */ protected Integer resolveReplicaCount(P primary, Context

context, AutoscalingSpec autoscaling, int staticReplicas, int initialReplicas) { + Optional existing = getSecondaryResource(primary, context); + Integer resolved = computeReplicaCount(primary, existing, autoscaling, + staticReplicas, initialReplicas); + logReplicaChange(primary, existing, resolved); + return resolved; + } + + private Integer computeReplicaCount(P primary, Optional existing, + AutoscalingSpec autoscaling, int staticReplicas, int initialReplicas) { // Suspended cluster → 0 replicas (dependent resources natively respect suspend). // Exception: HMS stays running if includeMetastore=false in autoSuspend config. if (primary instanceof HiveCluster hc && hc.getSpec().suspend()) { @@ -186,7 +196,6 @@ protected Integer resolveReplicaCount(P primary, Context

context, if (autoscaling == null || !autoscaling.isEnabled()) { return staticReplicas; } - Optional existing = getSecondaryResource(primary, context); if (existing.isPresent()) { // Check if the autoscaler has made a decision during this operator's lifecycle Integer managed = HiveClusterAutoscaler.getManagedReplicas( @@ -197,21 +206,37 @@ protected Integer resolveReplicaCount(P primary, Context

context, 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; + Integer current = Workloads.replicas(existing.get()); + return current != null ? current : initialReplicas; } // First creation: start at minReplicas. return initialReplicas; } + /** + * Emits an INFO line when the SSA about to run will actually change the workload's replica + * count. Silence means the count already matches, so no scale is happening. Without this, + * every scale of an HS2/Metastore Deployment reached the cluster silently (only the imperative + * LLAP path and the autoscaler logged); a bare "Reconciled" line said nothing about the size. + * Uses the same "Scaling ... A -> B" shape the imperative LLAP path emits. + */ + private void logReplicaChange(P primary, Optional existing, Integer desired) { + String component = getComponentName(); + if (component == null || desired == null) { + return; + } + String ns = primary.getMetadata().getNamespace(); + String name = existing.map(r -> r.getMetadata().getName()).orElse(component); + if (existing.isEmpty()) { + LOG.info("Creating {} {}/{} with {} replicas", component, ns, name, desired); + return; + } + Integer current = Workloads.replicas(existing.get()); + if (current != null && !current.equals(desired)) { + LOG.info("Scaling {} {}/{}: {} -> {} replicas", component, ns, name, current, desired); + } + } + /** * Returns the component name for this dependent (used for autoscaler replica lookup). diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java index 557ebe99d545..792820759ad8 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java @@ -60,6 +60,7 @@ import org.apache.hive.kubernetes.operator.model.status.ComponentStatus; import org.apache.hive.kubernetes.operator.util.ConfigUtils; import org.apache.hive.kubernetes.operator.util.Labels; +import org.apache.hive.kubernetes.operator.util.Workloads; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -617,6 +618,30 @@ private void patchReplicas(KubernetesClient client, HiveCluster resource, } } + /** + * Emits an INFO log when the reconciler's server-side apply is about to change the workload's + * replica count. Without this, an SSA-driven scale (a user editing spec.llapClusters[i].replicas, + * a helm upgrade rewriting it) reaches the StatefulSet/Deployment silently -- only the + * autoscaler path {@link #patchReplicas} logged its scales, so a plain scale looked like the + * operator was doing nothing. Read failures are swallowed at DEBUG: the SSA below runs either + * way, and a missing pre-scale line is not worth failing the reconcile over. + */ + private void logReplicaChange(KubernetesClient client, String ns, String workloadName, + String kind, int desired, boolean isStatefulSet) { + try { + Integer current = Workloads.replicas(isStatefulSet + ? client.apps().statefulSets().inNamespace(ns).withName(workloadName).get() + : client.apps().deployments().inNamespace(ns).withName(workloadName).get()); + if (current == null) { + LOG.info("Creating {} {}/{} with {} replicas", kind, ns, workloadName, desired); + } else if (current != desired) { + LOG.info("Scaling {} {}/{}: {} -> {} replicas", kind, ns, workloadName, current, desired); + } + } catch (Exception e) { + LOG.debug("Could not read current replicas for {}/{}: {}", ns, workloadName, e.getMessage()); + } + } + private void patchSuspendSpec(KubernetesClient client, HiveCluster resource, boolean suspend) { String ns = resource.getMetadata().getNamespace(); String name = resource.getMetadata().getName(); @@ -668,6 +693,8 @@ private void reconcileLlapClusters(HiveCluster resource, KubernetesClient client // brief scale-up-then-down on first create (K8s defaults to 1 if omitted). // resolveLlapReplicaCount already reads the autoscaler's managed value, // so this is always the correct replica count. + String llapWorkload = clusterName + "-" + llapSpec.name(); + logReplicaChange(client, ns, llapWorkload, "llap", replicas, /*isStatefulSet=*/true); client.apps().statefulSets().inNamespace(ns) .resource(LlapResourceBuilder.buildStatefulSet(resource, llapSpec, replicas)) .forceConflicts() @@ -687,6 +714,8 @@ private void reconcileLlapClusters(HiveCluster resource, KubernetesClient client client.services().inNamespace(ns) .resource(LlapResourceBuilder.buildTezAmService(resource, llapSpec)) .serverSideApply(); + String tezAmWorkload = clusterName + "-tezam-" + llapSpec.name(); + logReplicaChange(client, ns, tezAmWorkload, "tezam", tezAmReplicas, /*isStatefulSet=*/false); client.apps().deployments().inNamespace(ns) .resource(LlapResourceBuilder.buildTezAmDeployment(resource, llapSpec, tezAmReplicas)) .forceConflicts() diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java new file mode 100644 index 000000000000..973076894967 --- /dev/null +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java @@ -0,0 +1,47 @@ +/* + * 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 io.fabric8.kubernetes.api.model.HasMetadata; +import io.fabric8.kubernetes.api.model.apps.Deployment; +import io.fabric8.kubernetes.api.model.apps.StatefulSet; + +/** + * Small helpers over Deployment/StatefulSet that read fields without the + * null-guard boilerplate every caller would otherwise repeat. + */ +public final class Workloads { + + private Workloads() {} + + /** + * Returns spec.replicas from a Deployment or StatefulSet, or null when the resource is + * absent, has no spec, or the field is unset. A non-workload resource returns null too. + */ + public static Integer replicas(HasMetadata resource) { + if (resource instanceof Deployment d) { + return d.getSpec() == null ? null : d.getSpec().getReplicas(); + } + if (resource instanceof StatefulSet s) { + return s.getSpec() == null ? null : s.getSpec().getReplicas(); + } + return null; + } +} From e07c59eee139544e194e9402f33ade76666df85f Mon Sep 17 00:00:00 2001 From: Laszlo Bodor Date: Wed, 16 Sep 2026 11:16:43 +0200 Subject: [PATCH 2/3] refactor, Workloads improvement --- packaging/src/kubernetes/pom.xml | 31 +++++ .../dependent/LlapResourceBuilder.java | 15 +- .../reconciler/HiveClusterReconciler.java | 19 +-- .../kubernetes/operator/util/Workloads.java | 31 ++++- .../operator/util/TestWorkloads.java | 129 ++++++++++++++++++ 5 files changed, 203 insertions(+), 22 deletions(-) create mode 100644 packaging/src/kubernetes/src/test/java/org/apache/hive/kubernetes/operator/util/TestWorkloads.java diff --git a/packaging/src/kubernetes/pom.xml b/packaging/src/kubernetes/pom.xml index f1d8bbb84e5c..999d5ebcf0dc 100644 --- a/packaging/src/kubernetes/pom.xml +++ b/packaging/src/kubernetes/pom.xml @@ -93,6 +93,37 @@ log4j-core ${log4j2.version} + + + org.junit.jupiter + junit-jupiter-api + ${junit.jupiter.version} + test + + + org.junit.jupiter + junit-jupiter-engine + ${junit.jupiter.version} + test + + + org.junit.jupiter + junit-jupiter-params + ${junit.jupiter.version} + test + + + org.mockito + mockito-core + ${mockito-core.version} + test + + + org.mockito + mockito-junit-jupiter + ${mockito-core.version} + test + src/java 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 { private static final LlapResourceBuilder INSTANCE = new LlapResourceBuilder(); - private static final String TEZAM_INFIX = "-tezam-"; private static final String HIVE_CONFIG_VOLUME = "hive-config"; private static final String LLAP_CONFIG_VOLUME = "llap-config"; @@ -108,7 +108,7 @@ private static OwnerReference ownerRef(HiveCluster hc) { /** Resource name for a specific LLAP cluster: {clusterName}-{llapName}. */ public static String resourceName(HiveCluster hc, LlapSpec llap) { - return hc.getMetadata().getName() + "-" + llap.name(); + return Workloads.nameFor(hc, ConfigUtils.llapComponentKey(llap.name())); } /** ConfigMap name for a specific LLAP cluster. */ @@ -208,17 +208,22 @@ public static PodDisruptionBudget buildPdb(HiveCluster hc, LlapSpec llap) { /** TezAM Deployment/Service name for a specific LLAP cluster. */ public static String tezAmResourceName(HiveCluster hc, LlapSpec llap) { - return hc.getMetadata().getName() + TEZAM_INFIX + llap.name(); + return tezAmResourceName(hc, llap.name()); + } + + /** TezAM Deployment/Service name from an LLAP cluster name (used where only the name is in hand). */ + public static String tezAmResourceName(HiveCluster hc, String llapName) { + return Workloads.nameFor(hc, ConfigUtils.tezAmComponentKey(llapName)); } /** TezAM ConfigMap name for a specific LLAP cluster. */ public static String tezAmConfigMapName(HiveCluster hc, LlapSpec llap) { - return hc.getMetadata().getName() + TEZAM_INFIX + llap.name() + "-config"; + return tezAmResourceName(hc, llap) + "-config"; } /** TezAM PDB name for a specific LLAP cluster. */ public static String tezAmPdbName(HiveCluster hc, LlapSpec llap) { - return hc.getMetadata().getName() + TEZAM_INFIX + llap.name() + "-pdb"; + return tezAmResourceName(hc, llap) + "-pdb"; } /** Builds the PodDisruptionBudget for a per-LLAP-cluster TezAM. */ diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java index 792820759ad8..9be55af4cc4e 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java @@ -594,18 +594,7 @@ private static int getMinScrapeInterval(HiveClusterSpec spec) { private void patchReplicas(KubernetesClient client, HiveCluster resource, String component, int replicas) { String namespace = resource.getMetadata().getNamespace(); - // Component keys use prefixes: "llap-{name}" → workload "{cluster}-{name}", - // "tezam-{name}" → workload "{cluster}-tezam-{name}". - String workloadName; - if (component.startsWith(ConfigUtils.COMPONENT_LLAP + "-")) { - String llapName = component.substring(ConfigUtils.COMPONENT_LLAP.length() + 1); - workloadName = resource.getMetadata().getName() + "-" + llapName; - } else if (component.startsWith(ConfigUtils.COMPONENT_TEZAM + "-")) { - String llapName = component.substring(ConfigUtils.COMPONENT_TEZAM.length() + 1); - workloadName = resource.getMetadata().getName() + "-tezam-" + llapName; - } else { - workloadName = resource.getMetadata().getName() + "-" + component; - } + String workloadName = Workloads.nameFor(resource, component); try { if (component.startsWith(ConfigUtils.COMPONENT_LLAP + "-")) { client.apps().statefulSets().inNamespace(namespace).withName(workloadName).scale(replicas); @@ -714,7 +703,7 @@ private void reconcileLlapClusters(HiveCluster resource, KubernetesClient client client.services().inNamespace(ns) .resource(LlapResourceBuilder.buildTezAmService(resource, llapSpec)) .serverSideApply(); - String tezAmWorkload = clusterName + "-tezam-" + llapSpec.name(); + String tezAmWorkload = LlapResourceBuilder.tezAmResourceName(resource, llapSpec); logReplicaChange(client, ns, tezAmWorkload, "tezam", tezAmReplicas, /*isStatefulSet=*/false); client.apps().deployments().inNamespace(ns) .resource(LlapResourceBuilder.buildTezAmDeployment(resource, llapSpec, tezAmReplicas)) @@ -939,8 +928,8 @@ private boolean isClusterIdle(HiveCluster resource, KubernetesClient client) { if (spec.tezAm().isEnabled()) { for (var llap : spec.llapClusters()) { if (llap.isEnabled() - && !isAtMinReplicas(client, ns, name + "-tezam-" + llap.name(), false, - llap.tezAm().autoscaling().minReplicas())) { + && !isAtMinReplicas(client, ns, LlapResourceBuilder.tezAmResourceName(resource, llap), + false, llap.tezAm().autoscaling().minReplicas())) { return false; } } diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java index 973076894967..327b9198e9ae 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java @@ -22,10 +22,12 @@ import io.fabric8.kubernetes.api.model.HasMetadata; import io.fabric8.kubernetes.api.model.apps.Deployment; import io.fabric8.kubernetes.api.model.apps.StatefulSet; +import org.apache.hive.kubernetes.operator.model.HiveCluster; /** - * Small helpers over Deployment/StatefulSet that read fields without the - * null-guard boilerplate every caller would otherwise repeat. + * Helpers for the operator's workloads (Deployment/StatefulSet): reading fields without the + * null-guard boilerplate every caller would otherwise repeat, and resolving the workload's + * K8s name from the autoscaler's component key. */ public final class Workloads { @@ -44,4 +46,29 @@ public static Integer replicas(HasMetadata resource) { } return null; } + + /** + * 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..ec83d3548cea --- /dev/null +++ b/packaging/src/kubernetes/src/test/java/org/apache/hive/kubernetes/operator/util/TestWorkloads.java @@ -0,0 +1,129 @@ +/* + * 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.assertNull; + +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; + +class TestWorkloads { + + @Test + void replicasReturnsValueFromDeployment() { + Deployment d = new DeploymentBuilder() + .withNewSpec().withReplicas(3).endSpec() + .build(); + assertEquals(3, Workloads.replicas(d)); + } + + @Test + void replicasReturnsValueFromStatefulSet() { + StatefulSet s = new StatefulSetBuilder() + .withNewSpec().withReplicas(5).endSpec() + .build(); + assertEquals(5, Workloads.replicas(s)); + } + + @Test + void replicasReturnsNullWhenDeploymentSpecMissing() { + // fabric8 Deployment with no spec set at all + Deployment d = new DeploymentBuilder().build(); + assertNull(Workloads.replicas(d)); + } + + @Test + void replicasReturnsNullWhenStatefulSetSpecMissing() { + StatefulSet s = new StatefulSetBuilder().build(); + assertNull(Workloads.replicas(s)); + } + + @Test + void replicasReturnsNullWhenDeploymentReplicasUnset() { + // spec present, but replicas field not set + Deployment d = new DeploymentBuilder().withNewSpec().endSpec().build(); + assertNull(Workloads.replicas(d)); + } + + @Test + void replicasReturnsNullWhenStatefulSetReplicasUnset() { + StatefulSet s = new StatefulSetBuilder().withNewSpec().endSpec().build(); + assertNull(Workloads.replicas(s)); + } + + @Test + void replicasReturnsNullForNonWorkloadResource() { + // Anything that's neither a Deployment nor a StatefulSet returns null, + // even if it happens to have a "spec" (e.g., a ConfigMap here has none). + ConfigMap cm = new ConfigMapBuilder() + .withNewMetadata().withName("cm").endMetadata() + .build(); + assertNull(Workloads.replicas(cm)); + } + + @Test + void replicasReturnsNullForNullResource() { + assertNull(Workloads.replicas(null)); + } + + /** + * 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; + } +} From 5accd76f6801c3465d9ebdfe3fdc861e74da8e57 Mon Sep 17 00:00:00 2001 From: Laszlo Bodor Date: Mon, 21 Sep 2026 14:39:02 +0200 Subject: [PATCH 3/3] Address PR review comments - Join the "creating" and "scaling" log lines into one shared Workloads.logReplicaChange, used by both the dependents' SSA and the reconciler's imperative LLAP/TezAM SSAs. First creation renders as "none -> N" instead of a separate message. - Workloads.replicas returns OptionalInt so the absent case is in the type rather than a nullable Integer; callers use orElse(initialReplicas). - resolveReplicaCount returns int, making the documented "never null" guarantee compiler-enforced, and drops the dead null guards on getComponentName()/desired (both dependents that call it override getComponentName). - Log the exact workload name on create too, via getSecondaryResourceName() instead of falling back to the component name. - Use LlapResourceBuilder.resourceName()/ConfigUtils component constants at the reconciler call sites instead of re-deriving the name inline. --- .../dependent/HiveDependentResource.java | 51 +++++-------- .../HiveServer2DeploymentDependent.java | 2 +- .../MetastoreDeploymentDependent.java | 2 +- .../reconciler/HiveClusterReconciler.java | 33 +++++---- .../kubernetes/operator/util/Workloads.java | 44 +++++++++--- .../operator/util/TestWorkloads.java | 72 ++++++++++++++----- 6 files changed, 124 insertions(+), 80 deletions(-) diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java index daad1504747a..7e5944f0564f 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/HiveDependentResource.java @@ -162,10 +162,11 @@ protected R handleCreate(R desired, P primary, Context

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) @@ -174,16 +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 existing = getSecondaryResource(primary, context); - Integer resolved = computeReplicaCount(primary, existing, autoscaling, + int resolved = computeReplicaCount(primary, existing, autoscaling, staticReplicas, initialReplicas); - logReplicaChange(primary, existing, resolved); + // Without this, every scale of an HS2/Metastore Deployment reached the cluster silently: + // only the imperative LLAP path and the autoscaler logged, and a bare "Reconciled" line + // said nothing about the size. + Workloads.logReplicaChange(LOG, getComponentName(), primary.getMetadata().getNamespace(), + getSecondaryResourceName(primary, context), existing.orElse(null), resolved); return resolved; } - private Integer computeReplicaCount(P primary, Optional existing, + private int computeReplicaCount(P primary, Optional existing, AutoscalingSpec autoscaling, int staticReplicas, int initialReplicas) { // Suspended cluster → 0 replicas (dependent resources natively respect suspend). // Exception: HMS stays running if includeMetastore=false in autoSuspend config. @@ -205,39 +210,15 @@ private Integer computeReplicaCount(P primary, Optional existing, if (managed != null) { return managed; } - // Fallback: operator restarted and MANAGED_REPLICAS is empty — read current value - Integer current = Workloads.replicas(existing.get()); - return current != null ? current : 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; } - /** - * Emits an INFO line when the SSA about to run will actually change the workload's replica - * count. Silence means the count already matches, so no scale is happening. Without this, - * every scale of an HS2/Metastore Deployment reached the cluster silently (only the imperative - * LLAP path and the autoscaler logged); a bare "Reconciled" line said nothing about the size. - * Uses the same "Scaling ... A -> B" shape the imperative LLAP path emits. - */ - private void logReplicaChange(P primary, Optional existing, Integer desired) { - String component = getComponentName(); - if (component == null || desired == null) { - return; - } - String ns = primary.getMetadata().getNamespace(); - String name = existing.map(r -> r.getMetadata().getName()).orElse(component); - if (existing.isEmpty()) { - LOG.info("Creating {} {}/{} with {} replicas", component, ns, name, desired); - return; - } - Integer current = Workloads.replicas(existing.get()); - if (current != null && !current.equals(desired)) { - LOG.info("Scaling {} {}/{}: {} -> {} replicas", component, ns, name, current, desired); - } - } - - /** * 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/MetastoreDeploymentDependent.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/MetastoreDeploymentDependent.java index 73afedd9dbae..e72240899cf2 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/MetastoreDeploymentDependent.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/dependent/MetastoreDeploymentDependent.java @@ -142,7 +142,7 @@ protected Deployment desired(HiveCluster hiveCluster, AutoscalingSpec msAutoscaling = spec.metastore().autoscaling(); int initialReplicas = msAutoscaling != null && msAutoscaling.isEnabled() ? Math.max(1, msAutoscaling.minReplicas()) : spec.metastore().replicas(); - Integer replicas = resolveReplicaCount( + int replicas = resolveReplicaCount( hiveCluster, context, msAutoscaling, spec.metastore().replicas(), initialReplicas); Deployment deployment = new DeploymentBuilder() diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java index 9be55af4cc4e..c7f447ef0ac9 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/reconciler/HiveClusterReconciler.java @@ -608,24 +608,21 @@ private void patchReplicas(KubernetesClient client, HiveCluster resource, } /** - * Emits an INFO log when the reconciler's server-side apply is about to change the workload's - * replica count. Without this, an SSA-driven scale (a user editing spec.llapClusters[i].replicas, - * a helm upgrade rewriting it) reaches the StatefulSet/Deployment silently -- only the - * autoscaler path {@link #patchReplicas} logged its scales, so a plain scale looked like the - * operator was doing nothing. Read failures are swallowed at DEBUG: the SSA below runs either - * way, and a missing pre-scale line is not worth failing the reconcile over. + * Reads the workload's current replica count from the API server and hands it to + * {@link Workloads#logReplicaChange}, so an SSA-driven scale (a user editing + * spec.llapClusters[i].replicas, a helm upgrade rewriting it) logs the same line the dependents' + * SSA does instead of reaching the StatefulSet/Deployment silently. Only the read is local to + * this class; the message and the "log nothing when unchanged" rule are shared. Read failures + * are swallowed at DEBUG: the SSA below runs either way, and a missing pre-scale line is not + * worth failing the reconcile over. */ private void logReplicaChange(KubernetesClient client, String ns, String workloadName, - String kind, int desired, boolean isStatefulSet) { + String component, int desired, boolean isStatefulSet) { try { - Integer current = Workloads.replicas(isStatefulSet + HasMetadata current = isStatefulSet ? client.apps().statefulSets().inNamespace(ns).withName(workloadName).get() - : client.apps().deployments().inNamespace(ns).withName(workloadName).get()); - if (current == null) { - LOG.info("Creating {} {}/{} with {} replicas", kind, ns, workloadName, desired); - } else if (current != desired) { - LOG.info("Scaling {} {}/{}: {} -> {} replicas", kind, ns, workloadName, current, desired); - } + : client.apps().deployments().inNamespace(ns).withName(workloadName).get(); + Workloads.logReplicaChange(LOG, component, ns, workloadName, current, desired); } catch (Exception e) { LOG.debug("Could not read current replicas for {}/{}: {}", ns, workloadName, e.getMessage()); } @@ -682,8 +679,9 @@ private void reconcileLlapClusters(HiveCluster resource, KubernetesClient client // brief scale-up-then-down on first create (K8s defaults to 1 if omitted). // resolveLlapReplicaCount already reads the autoscaler's managed value, // so this is always the correct replica count. - String llapWorkload = clusterName + "-" + llapSpec.name(); - logReplicaChange(client, ns, llapWorkload, "llap", replicas, /*isStatefulSet=*/true); + String llapWorkload = LlapResourceBuilder.resourceName(resource, llapSpec); + logReplicaChange(client, ns, llapWorkload, ConfigUtils.COMPONENT_LLAP, replicas, + /*isStatefulSet=*/true); client.apps().statefulSets().inNamespace(ns) .resource(LlapResourceBuilder.buildStatefulSet(resource, llapSpec, replicas)) .forceConflicts() @@ -704,7 +702,8 @@ private void reconcileLlapClusters(HiveCluster resource, KubernetesClient client .resource(LlapResourceBuilder.buildTezAmService(resource, llapSpec)) .serverSideApply(); String tezAmWorkload = LlapResourceBuilder.tezAmResourceName(resource, llapSpec); - logReplicaChange(client, ns, tezAmWorkload, "tezam", tezAmReplicas, /*isStatefulSet=*/false); + logReplicaChange(client, ns, tezAmWorkload, ConfigUtils.COMPONENT_TEZAM, tezAmReplicas, + /*isStatefulSet=*/false); client.apps().deployments().inNamespace(ns) .resource(LlapResourceBuilder.buildTezAmDeployment(resource, llapSpec, tezAmReplicas)) .forceConflicts() diff --git a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java index 327b9198e9ae..d098cf8d6c07 100644 --- a/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java +++ b/packaging/src/kubernetes/src/java/org/apache/hive/kubernetes/operator/util/Workloads.java @@ -19,32 +19,56 @@ package org.apache.hive.kubernetes.operator.util; +import java.util.OptionalInt; + import io.fabric8.kubernetes.api.model.HasMetadata; import io.fabric8.kubernetes.api.model.apps.Deployment; import io.fabric8.kubernetes.api.model.apps.StatefulSet; import org.apache.hive.kubernetes.operator.model.HiveCluster; +import org.slf4j.Logger; /** * Helpers for the operator's workloads (Deployment/StatefulSet): reading fields without the - * null-guard boilerplate every caller would otherwise repeat, and resolving the workload's - * K8s name from the autoscaler's component key. + * null-guard boilerplate every caller would otherwise repeat, logging replica changes in one + * shape, and resolving the workload's K8s name from the autoscaler's component key. */ public final class Workloads { private Workloads() {} /** - * Returns spec.replicas from a Deployment or StatefulSet, or null when the resource is - * absent, has no spec, or the field is unset. A non-workload resource returns null too. + * Returns spec.replicas from a Deployment or StatefulSet. Empty when the resource is absent, + * has no spec, or the field is unset — in practice that means "the workload isn't there yet", + * since the API server defaults spec.replicas on write. A non-workload resource is empty too. */ - public static Integer replicas(HasMetadata resource) { - if (resource instanceof Deployment d) { - return d.getSpec() == null ? null : d.getSpec().getReplicas(); + public static OptionalInt replicas(HasMetadata resource) { + Integer replicas = null; + if (resource instanceof Deployment d && d.getSpec() != null) { + replicas = d.getSpec().getReplicas(); + } else if (resource instanceof StatefulSet s && s.getSpec() != null) { + replicas = s.getSpec().getReplicas(); } - if (resource instanceof StatefulSet s) { - return s.getSpec() == null ? null : s.getSpec().getReplicas(); + return replicas == null ? OptionalInt.empty() : OptionalInt.of(replicas); + } + + /** + * Logs the replica count the operator is about to apply to {@code namespace/name}. One line + * covers both cases: {@code current} is the workload as it exists in the cluster, or null when + * it doesn't exist yet, which logs as {@code none -> N}. Nothing is logged when the count + * already matches, so silence means no scale is happening. + *

+ * 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; } - return null; + log.info("Setting replica count for {} {}/{}: {} -> {}", component, namespace, name, + actual.isPresent() ? String.valueOf(actual.getAsInt()) : "none", desired); } /** 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 index ec83d3548cea..3d2e6d537796 100644 --- 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 @@ -20,7 +20,14 @@ package org.apache.hive.kubernetes.operator.util; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertNull; +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; @@ -33,6 +40,7 @@ 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 { @@ -41,7 +49,7 @@ void replicasReturnsValueFromDeployment() { Deployment d = new DeploymentBuilder() .withNewSpec().withReplicas(3).endSpec() .build(); - assertEquals(3, Workloads.replicas(d)); + assertEquals(OptionalInt.of(3), Workloads.replicas(d)); } @Test @@ -49,48 +57,80 @@ void replicasReturnsValueFromStatefulSet() { StatefulSet s = new StatefulSetBuilder() .withNewSpec().withReplicas(5).endSpec() .build(); - assertEquals(5, Workloads.replicas(s)); + assertEquals(OptionalInt.of(5), Workloads.replicas(s)); } @Test - void replicasReturnsNullWhenDeploymentSpecMissing() { + void replicasIsEmptyWhenDeploymentSpecMissing() { // fabric8 Deployment with no spec set at all Deployment d = new DeploymentBuilder().build(); - assertNull(Workloads.replicas(d)); + assertFalse(Workloads.replicas(d).isPresent()); } @Test - void replicasReturnsNullWhenStatefulSetSpecMissing() { + void replicasIsEmptyWhenStatefulSetSpecMissing() { StatefulSet s = new StatefulSetBuilder().build(); - assertNull(Workloads.replicas(s)); + assertFalse(Workloads.replicas(s).isPresent()); } @Test - void replicasReturnsNullWhenDeploymentReplicasUnset() { + void replicasIsEmptyWhenDeploymentReplicasUnset() { // spec present, but replicas field not set Deployment d = new DeploymentBuilder().withNewSpec().endSpec().build(); - assertNull(Workloads.replicas(d)); + assertFalse(Workloads.replicas(d).isPresent()); } @Test - void replicasReturnsNullWhenStatefulSetReplicasUnset() { + void replicasIsEmptyWhenStatefulSetReplicasUnset() { StatefulSet s = new StatefulSetBuilder().withNewSpec().endSpec().build(); - assertNull(Workloads.replicas(s)); + assertFalse(Workloads.replicas(s).isPresent()); } @Test - void replicasReturnsNullForNonWorkloadResource() { - // Anything that's neither a Deployment nor a StatefulSet returns null, + 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(); - assertNull(Workloads.replicas(cm)); + 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 replicasReturnsNullForNullResource() { - assertNull(Workloads.replicas(null)); + 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); } /**