diff --git a/llap-server/src/java/org/apache/hadoop/hive/llap/cli/status/LlapStatusServiceDriver.java b/llap-server/src/java/org/apache/hadoop/hive/llap/cli/status/LlapStatusServiceDriver.java index a16e18f85fca..88fe6c8091ab 100644 --- a/llap-server/src/java/org/apache/hadoop/hive/llap/cli/status/LlapStatusServiceDriver.java +++ b/llap-server/src/java/org/apache/hadoop/hive/llap/cli/status/LlapStatusServiceDriver.java @@ -71,7 +71,7 @@ public class LlapStatusServiceDriver { private static final Logger CONSOLE_LOGGER = LoggerFactory.getLogger("LlapStatusServiceDriverConsole"); private static final EnumSet NO_YARN_SERVICE_INFO_STATES = EnumSet.of( - State.APP_NOT_FOUND, State.COMPLETE, State.LAUNCHING); + State.COMPLETE, State.LAUNCHING); private static final EnumSet LAUNCHING_STATES = EnumSet.of( State.LAUNCHING, State.RUNNING_PARTIAL, State.RUNNING_ALL); @@ -107,6 +107,10 @@ public class LlapStatusServiceDriver { private static final long LOG_SUMMARY_INTERVAL = 15000L; // Log summary every ~15 seconds. private static final String LLAP_KEY = "llap"; + private static final String TEZ_AM_FRAMEWORK_MODE = "tez.am.framework.mode"; + private static final String TEZ_FRAMEWORK_MODE_STANDALONE_ZOOKEEPER = "STANDALONE_ZOOKEEPER"; + /** When set, overrides auto-detection for registry vs YARN status lookup. */ + private static final String CONFIG_STATUS_USE_REGISTRY = CONF_PREFIX + "status.use-registry"; private final Configuration conf; private String appName = null; @@ -194,18 +198,31 @@ public ExitCode run(LlapStatusServiceCommandLine cl, long watchTimeoutMs) { llapRegistryConf.set(HiveConf.ConfVars.LLAP_DAEMON_SERVICE_HOSTS.varname, "@" + appName); } + if (usesRegistryBasedLlapStatus(conf)) { + LOG.info("Non-YARN LLAP deployment detected; using LLAP registry for status"); + try { + ExitCode ret = populateAppStatusFromLlapRegistry(appStatusBuilder, watchTimeoutMs, true); + if (ret == ExitCode.SUCCESS) { + updateRunningThresholdAchieved(appStatusBuilder, cl.getRunningNodesThreshold()); + } + return ret; + } catch (LlapStatusCliException e) { + logError(e); + return e.getExitCode(); + } + } + try { if (serviceClient == null) { serviceClient = LlapSliderUtils.createServiceClient(conf); } } catch (Exception e) { - LlapStatusCliException le = new LlapStatusCliException( - ExitCode.SERVICE_CLIENT_ERROR_CREATE_FAILED, "Failed to create service client", e); + LlapStatusCliException le = new LlapStatusCliException(ExitCode.SERVICE_CLIENT_ERROR_CREATE_FAILED, + "Failed to create YARN Service client", e); logError(le); return le.getExitCode(); } - // Get the App report from YARN ApplicationReport appReport; try { appReport = getAppReport(appName, cl.getFindAppTimeoutMs()); @@ -214,7 +231,6 @@ public ExitCode run(LlapStatusServiceCommandLine cl, long watchTimeoutMs) { return e.getExitCode(); } - // Process the report ExitCode ret; try { ret = processAppReport(appReport, appStatusBuilder); @@ -225,30 +241,35 @@ public ExitCode run(LlapStatusServiceCommandLine cl, long watchTimeoutMs) { if (ret != ExitCode.SUCCESS) { return ret; - } else if (NO_YARN_SERVICE_INFO_STATES.contains(appStatusBuilder.getState())) { + } + + if (NO_YARN_SERVICE_INFO_STATES.contains(appStatusBuilder.getState())) { + updateRunningThresholdAchieved(appStatusBuilder, cl.getRunningNodesThreshold()); return ExitCode.SUCCESS; - } else { - // Get information from YARN Service - try { - ret = populateAppStatusFromServiceStatus(appName, serviceClient, appStatusBuilder); - } catch (LlapStatusCliException e) { - // In case of failure, send back whatever is constructed so far - which would be from the AppReport - logError(e); - return e.getExitCode(); - } + } + + // Get information from YARN Service + try { + ret = populateAppStatusFromServiceStatus(appName, serviceClient, appStatusBuilder); + } catch (LlapStatusCliException e) { + // In case of failure, send back whatever is constructed so far - which would be from the AppReport + logError(e); + return e.getExitCode(); } if (ret != ExitCode.SUCCESS) { return ret; - } else { - try { - ret = populateAppStatusFromLlapRegistry(appStatusBuilder, watchTimeoutMs); - } catch (LlapStatusCliException e) { - logError(e); - return e.getExitCode(); - } } + try { + ret = populateAppStatusFromLlapRegistry(appStatusBuilder, watchTimeoutMs, false); + } catch (LlapStatusCliException e) { + logError(e); + return e.getExitCode(); + } + if (ret == ExitCode.SUCCESS) { + updateRunningThresholdAchieved(appStatusBuilder, cl.getRunningNodesThreshold()); + } return ret; } finally { LOG.debug("Final AppState: " + appStatusBuilder.toString()); @@ -404,6 +425,11 @@ private ExitCode populateAppStatusFromServiceStatus(String appName, ServiceClien */ private ExitCode populateAppStatusFromLlapRegistry(AppStatusBuilder appStatusBuilder, long watchTimeoutMs) throws LlapStatusCliException { + return populateAppStatusFromLlapRegistry(appStatusBuilder, watchTimeoutMs, false); + } + + private ExitCode populateAppStatusFromLlapRegistry(AppStatusBuilder appStatusBuilder, long watchTimeoutMs, + boolean registryOnly) throws LlapStatusCliException { if (llapRegistry == null) { try { @@ -427,6 +453,23 @@ private ExitCode populateAppStatusFromLlapRegistry(AppStatusBuilder appStatusBui appStatusBuilder.setState(State.LAUNCHING); appStatusBuilder.clearRunningLlapInstances(); return ExitCode.SUCCESS; + } + + if (registryOnly) { + List registryInstances = new LinkedList<>(); + for (LlapServiceInstance serviceInstance : serviceInstances) { + registryInstances.add(createLlapInstanceFromRegistry(serviceInstance)); + } + if (appStatusBuilder.getAmInfo() == null) { + appStatusBuilder.setAmInfo(new AmInfo().setAppName(appName)); + } + appStatusBuilder.clearAndAddPreviouslyKnownRunningInstances(registryInstances); + appStatusBuilder.setLiveInstances(registryInstances.size()); + if (appStatusBuilder.getDesiredInstances() == null) { + appStatusBuilder.setDesiredInstances(registryInstances.size()); + } + updateStateFromInstanceCounts(appStatusBuilder, registryInstances.size()); + return ExitCode.SUCCESS; } else { // Tracks instances known by both YARN Service and llap. List validatedInstances = new LinkedList<>(); @@ -455,18 +498,10 @@ private ExitCode populateAppStatusFromLlapRegistry(AppStatusBuilder appStatusBui appStatusBuilder.setLiveInstances(validatedInstances.size()); appStatusBuilder.setLaunchingInstances(llapExtraInstances.size()); - if (appStatusBuilder.getDesiredInstances() != null && - validatedInstances.size() >= appStatusBuilder.getDesiredInstances()) { - appStatusBuilder.setState(State.RUNNING_ALL); - if (validatedInstances.size() > appStatusBuilder.getDesiredInstances()) { - LOG.warn("Found more entries in LLAP registry, as compared to desired entries"); - } - } else { - if (validatedInstances.size() > 0) { - appStatusBuilder.setState(State.RUNNING_PARTIAL); - } else { - appStatusBuilder.setState(State.LAUNCHING); - } + updateStateFromInstanceCounts(appStatusBuilder, validatedInstances.size()); + if (appStatusBuilder.getDesiredInstances() != null + && validatedInstances.size() > appStatusBuilder.getDesiredInstances()) { + LOG.warn("Found more entries in LLAP registry, as compared to desired entries"); } // At this point, everything that can be consumed from AppStatusBuilder has been consumed. @@ -487,6 +522,63 @@ private ExitCode populateAppStatusFromLlapRegistry(AppStatusBuilder appStatusBui return ExitCode.SUCCESS; } + static LlapInstance createLlapInstanceFromRegistry(LlapServiceInstance serviceInstance) { + String containerId = serviceInstance.getProperties().get( + HiveConf.ConfVars.LLAP_DAEMON_CONTAINER_ID.varname); + if (StringUtils.isBlank(containerId)) { + containerId = serviceInstance.getWorkerIdentity(); + } + LlapInstance llapInstance = new LlapInstance(serviceInstance.getHost(), containerId); + llapInstance.setMgmtPort(serviceInstance.getManagementPort()); + llapInstance.setRpcPort(serviceInstance.getRpcPort()); + llapInstance.setShufflePort(serviceInstance.getShufflePort()); + llapInstance.setWebUrl(serviceInstance.getServicesAddress()); + llapInstance.setStatusUrl(serviceInstance.getServicesAddress() + "/status"); + return llapInstance; + } + + static void updateStateFromInstanceCounts(AppStatusBuilder appStatusBuilder, int liveInstances) { + Integer desiredInstances = appStatusBuilder.getDesiredInstances(); + if (desiredInstances != null && liveInstances >= desiredInstances) { + appStatusBuilder.setState(State.RUNNING_ALL); + } else if (liveInstances > 0) { + appStatusBuilder.setState(State.RUNNING_PARTIAL); + } else { + appStatusBuilder.setState(State.LAUNCHING); + } + } + + /** + * Non-YARN LLAP (e.g. Kubernetes, standalone Tez) uses the ZK registry for daemon discovery. + * YARN Service LLAP uses the YARN APIs. Detection follows deployment config set by the operator. + */ + static boolean usesRegistryBasedLlapStatus(Configuration conf) { + if (conf.get(CONFIG_STATUS_USE_REGISTRY) != null) { + return conf.getBoolean(CONFIG_STATUS_USE_REGISTRY, false); + } + if (HiveConf.getBoolVar(conf, HiveConf.ConfVars.HIVE_SERVER2_TEZ_USE_EXTERNAL_SESSIONS)) { + return true; + } + String tezFrameworkMode = conf.get(TEZ_AM_FRAMEWORK_MODE, ""); + return TEZ_FRAMEWORK_MODE_STANDALONE_ZOOKEEPER.equalsIgnoreCase(tezFrameworkMode); + } + + static void updateRunningThresholdAchieved(AppStatusBuilder appStatusBuilder, float runningNodesThreshold) { + Integer desiredInstances = appStatusBuilder.getDesiredInstances(); + Integer liveInstances = appStatusBuilder.getLiveInstances(); + if (desiredInstances == null || desiredInstances <= 0 || liveInstances == null) { + return; + } + if (!(appStatusBuilder.getState() == State.RUNNING_PARTIAL + || appStatusBuilder.getState() == State.RUNNING_ALL)) { + return; + } + float ratio = (float) liveInstances / (float) desiredInstances; + if (ratio >= runningNodesThreshold) { + appStatusBuilder.setRunningThresholdAchieved(true); + } + } + private void close() { if (serviceClient != null) { serviceClient.stop(); diff --git a/llap-server/src/test/org/apache/hadoop/hive/llap/cli/status/TestLlapStatusRegistryFallback.java b/llap-server/src/test/org/apache/hadoop/hive/llap/cli/status/TestLlapStatusRegistryFallback.java new file mode 100644 index 000000000000..a1773b726b2b --- /dev/null +++ b/llap-server/src/test/org/apache/hadoop/hive/llap/cli/status/TestLlapStatusRegistryFallback.java @@ -0,0 +1,198 @@ +/* + * 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.hadoop.hive.llap.cli.status; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hive.conf.HiveConf; +import org.apache.hadoop.hive.llap.registry.LlapServiceInstance; +import org.apache.hadoop.yarn.api.records.Resource; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +/** + * Tests registry-only status helpers used when LLAP runs without YARN. + */ +public class TestLlapStatusRegistryFallback { + + @Test + public void testCreateLlapInstanceFromRegistryUsesWorkerIdentityWhenNoContainerId() { + LlapServiceInstance instance = new TestLlapServiceInstance("worker-1", "llap-host-0", + "http://llap-host-0:15002", 15004, 15001, 15551, Collections.emptyMap()); + + LlapInstance llapInstance = LlapStatusServiceDriver.createLlapInstanceFromRegistry(instance); + + assertEquals("llap-host-0", llapInstance.getHostname()); + assertEquals("worker-1", llapInstance.getContainerId()); + assertEquals("http://llap-host-0:15002", llapInstance.getWebUrl()); + assertEquals("http://llap-host-0:15002/status", llapInstance.getStatusUrl()); + assertEquals(Integer.valueOf(15004), llapInstance.getMgmtPort()); + } + + @Test + public void testCreateLlapInstanceFromRegistryUsesYarnContainerIdWhenPresent() { + Map props = new HashMap<>(); + props.put(HiveConf.ConfVars.LLAP_DAEMON_CONTAINER_ID.varname, "container_123_456"); + LlapServiceInstance instance = new TestLlapServiceInstance("worker-2", "yarn-host", + "http://yarn-host:15002", 15004, 15001, 15551, props); + + LlapInstance llapInstance = LlapStatusServiceDriver.createLlapInstanceFromRegistry(instance); + + assertEquals("container_123_456", llapInstance.getContainerId()); + } + + @Test + public void testUpdateStateFromInstanceCounts() { + AppStatusBuilder builder = new AppStatusBuilder(); + builder.setDesiredInstances(2); + + LlapStatusServiceDriver.updateStateFromInstanceCounts(builder, 2); + assertEquals(State.RUNNING_ALL, builder.getState()); + + LlapStatusServiceDriver.updateStateFromInstanceCounts(builder, 1); + assertEquals(State.RUNNING_PARTIAL, builder.getState()); + + LlapStatusServiceDriver.updateStateFromInstanceCounts(builder, 0); + assertEquals(State.LAUNCHING, builder.getState()); + } + + @Test + public void testUsesRegistryBasedLlapStatusFromExternalSessions() { + Configuration conf = new Configuration(false); + HiveConf.setBoolVar(conf, HiveConf.ConfVars.HIVE_SERVER2_TEZ_USE_EXTERNAL_SESSIONS, true); + assertTrue(LlapStatusServiceDriver.usesRegistryBasedLlapStatus(conf)); + } + + @Test + public void testUsesRegistryBasedLlapStatusFromTezFrameworkMode() { + Configuration conf = new Configuration(false); + conf.set("tez.am.framework.mode", "STANDALONE_ZOOKEEPER"); + assertTrue(LlapStatusServiceDriver.usesRegistryBasedLlapStatus(conf)); + } + + @Test + public void testUsesYarnBasedLlapStatusByDefault() { + Configuration conf = new Configuration(false); + assertFalse(LlapStatusServiceDriver.usesRegistryBasedLlapStatus(conf)); + } + + @Test + public void testUpdateRunningThresholdAchievedWhenFullyRunning() { + AppStatusBuilder builder = new AppStatusBuilder(); + builder.setDesiredInstances(2); + builder.setLiveInstances(2); + builder.setState(State.RUNNING_ALL); + + LlapStatusServiceDriver.updateRunningThresholdAchieved(builder, 1.0f); + assertTrue(builder.isRunningThresholdAchieved()); + } + + @Test + public void testUpdateRunningThresholdAchievedWhenLaunching() { + AppStatusBuilder builder = new AppStatusBuilder(); + builder.setDesiredInstances(2); + builder.setLiveInstances(0); + builder.setState(State.LAUNCHING); + + LlapStatusServiceDriver.updateRunningThresholdAchieved(builder, 1.0f); + assertFalse(builder.isRunningThresholdAchieved()); + } + + private static final class TestLlapServiceInstance implements LlapServiceInstance { + private final String workerIdentity; + private final String host; + private final String servicesAddress; + private final int mgmtPort; + private final int rpcPort; + private final int shufflePort; + private final Map properties; + + private TestLlapServiceInstance(String workerIdentity, String host, String servicesAddress, + int mgmtPort, int rpcPort, int shufflePort, Map properties) { + this.workerIdentity = workerIdentity; + this.host = host; + this.servicesAddress = servicesAddress; + this.mgmtPort = mgmtPort; + this.rpcPort = rpcPort; + this.shufflePort = shufflePort; + this.properties = properties; + } + + @Override + public String getWorkerIdentity() { + return workerIdentity; + } + + @Override + public String getHost() { + return host; + } + + @Override + public int getRpcPort() { + return rpcPort; + } + + @Override + public Map getProperties() { + return properties; + } + + @Override + public int getManagementPort() { + return mgmtPort; + } + + @Override + public int getShufflePort() { + return shufflePort; + } + + @Override + public String getServicesAddress() { + return servicesAddress; + } + + @Override + public int getOutputFormatPort() { + return 0; + } + + @Override + public String getExternalHostname() { + return host; + } + + @Override + public int getExternalClientsRpcPort() { + return 0; + } + + @Override + public Resource getResource() { + return null; + } + } +} diff --git a/service/src/java/org/apache/hive/http/LlapServlet.java b/service/src/java/org/apache/hive/http/LlapServlet.java index 340e075800ab..6a5e5e0ede86 100644 --- a/service/src/java/org/apache/hive/http/LlapServlet.java +++ b/service/src/java/org/apache/hive/http/LlapServlet.java @@ -100,6 +100,8 @@ public void doGet(HttpServletRequest request, HttpServletResponse response) { ExitCode ret = driver.run(LlapStatusServiceCommandLine.parseArguments(new String[] {"-n", clusterName}), 0); if (ret == ExitCode.SUCCESS) { driver.outputJson(writer); + } else { + response.setStatus(HttpServletResponse.SC_INTERNAL_SERVER_ERROR); } } finally { diff --git a/service/src/resources/hive-webapps/hiveserver2/llap.html b/service/src/resources/hive-webapps/hiveserver2/llap.html index 62ac1b09e417..bf5522b3f48e 100644 --- a/service/src/resources/hive-webapps/hiveserver2/llap.html +++ b/service/src/resources/hive-webapps/hiveserver2/llap.html @@ -15,8 +15,8 @@ - - + +