Skip to content

Commit fa844c0

Browse files
HIVE-27126: queue level resource stats for YARN RM.
1 parent b335b07 commit fa844c0

14 files changed

Lines changed: 1582 additions & 1 deletion

File tree

beeline/src/java/org/apache/hive/beeline/logs/BeelineInPlaceUpdateStream.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,5 +96,10 @@ public String executionStatus() {
9696
public double progressedPercentage() {
9797
return response.getProgressedPercentage();
9898
}
99+
100+
@Override
101+
public String queueMetrics() {
102+
return "";
103+
}
99104
}
100105
}

common/src/java/org/apache/hadoop/hive/common/log/InPlaceUpdate.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -175,6 +175,12 @@ public void render(ProgressMonitor monitor) {
175175
reprintLine(SEPARATOR);
176176
reprintLineWithColorAsBold(footer, Ansi.Color.RED);
177177
reprintLine(SEPARATOR);
178+
179+
// Display queue metrics if available (may be multi-line: queue name + metrics)
180+
String queueMetrics = monitor.queueMetrics();
181+
if (queueMetrics != null && !queueMetrics.isEmpty()) {
182+
reprintMultiLine(queueMetrics);
183+
}
178184
}
179185

180186

common/src/java/org/apache/hadoop/hive/common/log/ProgressMonitor.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,11 @@ public String executionStatus() {
5252
public double progressedPercentage() {
5353
return 0;
5454
}
55+
56+
@Override
57+
public String queueMetrics() {
58+
return "";
59+
}
5560
};
5661

5762
List<String> headers();
@@ -65,4 +70,6 @@ public double progressedPercentage() {
6570
String executionStatus();
6671

6772
double progressedPercentage();
73+
74+
String queueMetrics();
6875
}

common/src/java/org/apache/hadoop/hive/conf/HiveConf.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3982,6 +3982,12 @@ public static enum ConfVars {
39823982
HIVE_SERVER2_TEZ_QUEUE_ACCESS_CHECK("hive.server2.tez.queue.access.check", false,
39833983
"Whether to check user access to explicitly specified YARN queues. " +
39843984
"yarn.resourcemanager.webapp.address must be configured to use this."),
3985+
HIVE_TEZ_QUEUE_METRICS_REFRESH_INTERVAL("hive.tez.queue.metrics.refresh.interval", "0s",
3986+
new TimeValidator(TimeUnit.SECONDS),
3987+
"Interval for refreshing YARN queue resource metrics during Tez query execution. " +
3988+
"When set to a positive value (e.g. 10s), displays real-time memory, vCore, capacity " +
3989+
"and application metrics for the YARN queue being used. " +
3990+
"Set to 0 or negative to disable. Minimum effective value is 1 second."),
39853991
HIVE_SERVER2_TEZ_SESSION_LIFETIME("hive.server2.tez.session.lifetime", "162h",
39863992
new TimeValidator(TimeUnit.HOURS),
39873993
"The lifetime of the Tez sessions launched by HS2 when default sessions are enabled.\n" +

ql/src/java/org/apache/hadoop/hive/ql/exec/tez/TezSession.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333
import org.apache.hadoop.hive.ql.session.SessionState.LogHelper;
3434
import org.apache.hadoop.hive.ql.wm.WmContext;
3535
import org.apache.hadoop.yarn.api.records.LocalResource;
36+
import org.apache.hadoop.yarn.client.api.YarnClient;
3637
import org.apache.tez.client.TezClient;
3738
import org.apache.tez.dag.api.TezException;
3839
import org.apache.tez.dag.api.client.DAGStatus;
@@ -86,6 +87,7 @@ public String toString() {
8687

8788
HiveConf getConf();
8889
TezClient getTezClient();
90+
YarnClient getYarnClient();
8991
boolean isOpen();
9092
boolean isOpening();
9193
boolean getDoAsEnabled();

ql/src/java/org/apache/hadoop/hive/ql/exec/tez/TezSessionPoolSession.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@
3333
import org.apache.hadoop.hive.ql.wm.WmContext;
3434
import org.apache.hadoop.hive.registry.impl.TezAmInstance;
3535
import org.apache.hadoop.yarn.api.records.LocalResource;
36+
import org.apache.hadoop.yarn.client.api.YarnClient;
3637
import org.apache.tez.client.TezClient;
3738
import org.apache.tez.dag.api.TezException;
3839
import org.apache.tez.dag.api.client.DAGStatus;
@@ -337,6 +338,11 @@ public TezClient getTezClient() {
337338
return baseSession.getTezClient();
338339
}
339340

341+
@Override
342+
public YarnClient getYarnClient() {
343+
return baseSession.getYarnClient();
344+
}
345+
340346
@Override
341347
public boolean isOpening() {
342348
return baseSession.isOpening();

ql/src/java/org/apache/hadoop/hive/ql/exec/tez/TezSessionState.java

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@
7373
import org.apache.hadoop.security.token.Token;
7474
import org.apache.hadoop.yarn.api.records.LocalResource;
7575
import org.apache.hadoop.yarn.api.records.LocalResourceType;
76+
import org.apache.hadoop.yarn.client.api.YarnClient;
7677
import org.apache.hadoop.yarn.conf.YarnConfiguration;
7778
import org.apache.tez.client.TezClient;
7879
import org.apache.tez.common.TezUtils;
@@ -117,6 +118,7 @@ public class TezSessionState implements TezSession {
117118
Path tezScratchDir;
118119
protected LocalResource appJarLr;
119120
private TezClient session;
121+
private YarnClient yarnClient;
120122
private Future<TezClient> sessionFuture;
121123
/** Console used for user feedback during async session opening. */
122124
private LogHelper console;
@@ -750,6 +752,17 @@ public void close(boolean keepDagFilesDir) throws Exception {
750752
closeClient(asyncSession);
751753
}
752754
}
755+
756+
// Stop YarnClient if it was initialized
757+
if (yarnClient != null) {
758+
try {
759+
LOG.info("Stopping YarnClient for session: {}", sessionId);
760+
yarnClient.stop();
761+
yarnClient = null;
762+
} catch (Exception e) {
763+
LOG.warn("Error stopping YarnClient for session {}: {}", sessionId, e.getMessage());
764+
}
765+
}
753766
} finally {
754767
try {
755768
cleanupScratchDir();
@@ -795,6 +808,20 @@ public String getSessionId() {
795808

796809
protected final void setTezClient(TezClient session) {
797810
this.session = session;
811+
812+
// Initialize YarnClient for queue metrics collection
813+
if (session != null && yarnClient == null) {
814+
try {
815+
yarnClient = YarnClient.createYarnClient();
816+
yarnClient.init(conf);
817+
yarnClient.start();
818+
LOG.info("YarnClient initialized for session: {}", sessionId);
819+
} catch (Exception e) {
820+
LOG.warn("Failed to initialize YarnClient for metrics collection: {}", e.getMessage());
821+
LOG.debug("Full exception for YarnClient initialization failure", e);
822+
yarnClient = null;
823+
}
824+
}
798825
}
799826

800827
@Override
@@ -820,6 +847,11 @@ public TezClient getTezClient() {
820847
return session;
821848
}
822849

850+
@Override
851+
public YarnClient getYarnClient() {
852+
return yarnClient;
853+
}
854+
823855
@Override
824856
public LocalResource getAppJarLr() {
825857
return appJarLr;

0 commit comments

Comments
 (0)