FrankChen021 commented on code in PR #19855: URL: https://github.com/apache/druid/pull/19855#discussion_r3698986109
########## sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemStackTraceTable.java: ########## @@ -0,0 +1,378 @@ +/* + * 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.druid.sql.calcite.schema; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.base.Preconditions; +import org.apache.calcite.DataContext; +import org.apache.calcite.linq4j.Enumerable; +import org.apache.calcite.linq4j.Linq4j; +import org.apache.calcite.rel.type.RelDataType; +import org.apache.calcite.rel.type.RelDataTypeFactory; +import org.apache.calcite.rex.RexNode; +import org.apache.calcite.schema.ProjectableFilterableTable; +import org.apache.calcite.schema.Schema; +import org.apache.calcite.schema.impl.AbstractTable; +import org.apache.druid.discovery.DiscoveryDruidNode; +import org.apache.druid.discovery.DruidNodeDiscoveryProvider; +import org.apache.druid.error.InvalidInput; +import org.apache.druid.java.util.common.StringUtils; +import org.apache.druid.java.util.common.logger.Logger; +import org.apache.druid.java.util.http.client.HttpClient; +import org.apache.druid.java.util.http.client.Request; +import org.apache.druid.java.util.http.client.response.StringFullResponseHandler; +import org.apache.druid.java.util.http.client.response.StringFullResponseHolder; +import org.apache.druid.query.QueryContexts; +import org.apache.druid.segment.column.ColumnType; +import org.apache.druid.segment.column.RowSignature; +import org.apache.druid.server.DruidNode; +import org.apache.druid.server.StackTraceCollector; +import org.apache.druid.server.security.AuthenticationResult; +import org.apache.druid.server.security.AuthorizerMapper; +import org.apache.druid.sql.calcite.planner.PlannerContext; +import org.apache.druid.sql.calcite.table.RowSignatures; +import org.jboss.netty.handler.codec.http.HttpMethod; + +import javax.annotation.Nullable; +import javax.servlet.http.HttpServletResponse; +import java.net.URI; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.Iterator; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; + +/** + * System schema table {@code sys.stack_trace} that contains a live Java thread-stack snapshot for + * explicitly selected Druid servers. + */ +public class SystemStackTraceTable extends AbstractTable implements ProjectableFilterableTable +{ + private static final Logger log = new Logger(SystemStackTraceTable.class); + + public static final String TABLE_NAME = "stack_trace"; + + static final RowSignature ROW_SIGNATURE = RowSignature + .builder() + .add("server", ColumnType.STRING) + .add("service_name", ColumnType.STRING) + .add("node_roles", ColumnType.STRING) + .add("collected_at", ColumnType.STRING) + .add("thread_id", ColumnType.LONG) + .add("thread_name", ColumnType.STRING) + .add("thread_state", ColumnType.STRING) + .add("daemon", ColumnType.LONG) + .add("priority", ColumnType.LONG) + .add("cpu_time_ns", ColumnType.LONG) + .add("user_cpu_time_ns", ColumnType.LONG) + .add("lock_name", ColumnType.STRING) + .add("lock_owner_id", ColumnType.LONG) + .add("lock_owner_name", ColumnType.STRING) + .add("is_deadlocked", ColumnType.LONG) + .add("stack", ColumnType.STRING) + .add("error_message", ColumnType.STRING) + .build(); + + private static final int SERVER_INDEX = ROW_SIGNATURE.indexOf("server"); + private static final int SERVICE_NAME_INDEX = ROW_SIGNATURE.indexOf("service_name"); + private static final int NODE_ROLES_INDEX = ROW_SIGNATURE.indexOf("node_roles"); + private static final int COLLECTED_AT_INDEX = ROW_SIGNATURE.indexOf("collected_at"); + private static final int THREAD_ID_INDEX = ROW_SIGNATURE.indexOf("thread_id"); + private static final int THREAD_NAME_INDEX = ROW_SIGNATURE.indexOf("thread_name"); + private static final int THREAD_STATE_INDEX = ROW_SIGNATURE.indexOf("thread_state"); + private static final int DAEMON_INDEX = ROW_SIGNATURE.indexOf("daemon"); + private static final int PRIORITY_INDEX = ROW_SIGNATURE.indexOf("priority"); + private static final int CPU_TIME_NS_INDEX = ROW_SIGNATURE.indexOf("cpu_time_ns"); + private static final int USER_CPU_TIME_NS_INDEX = ROW_SIGNATURE.indexOf("user_cpu_time_ns"); + private static final int LOCK_NAME_INDEX = ROW_SIGNATURE.indexOf("lock_name"); + private static final int LOCK_OWNER_ID_INDEX = ROW_SIGNATURE.indexOf("lock_owner_id"); + private static final int LOCK_OWNER_NAME_INDEX = ROW_SIGNATURE.indexOf("lock_owner_name"); + private static final int IS_DEADLOCKED_INDEX = ROW_SIGNATURE.indexOf("is_deadlocked"); + private static final int STACK_INDEX = ROW_SIGNATURE.indexOf("stack"); + private static final int ERROR_MESSAGE_INDEX = ROW_SIGNATURE.indexOf("error_message"); + + private final DruidNodeDiscoveryProvider druidNodeDiscoveryProvider; + private final AuthorizerMapper authorizerMapper; + private final HttpClient httpClient; + private final ObjectMapper jsonMapper; + + public SystemStackTraceTable( + final DruidNodeDiscoveryProvider druidNodeDiscoveryProvider, + final AuthorizerMapper authorizerMapper, + final HttpClient httpClient, + final ObjectMapper jsonMapper + ) + { + this.druidNodeDiscoveryProvider = druidNodeDiscoveryProvider; + this.authorizerMapper = authorizerMapper; + this.httpClient = httpClient; + this.jsonMapper = jsonMapper; + } + + @Override + public RelDataType getRowType(final RelDataTypeFactory typeFactory) + { + return RowSignatures.toRelDataType(ROW_SIGNATURE, typeFactory); + } + + @Override + public Schema.TableType getJdbcTableType() + { + return Schema.TableType.SYSTEM_TABLE; + } + + @Override + public Enumerable<Object[]> scan( + final DataContext root, + final List<RexNode> filters, + @Nullable final int[] projects + ) + { + final AuthenticationResult authenticationResult = (AuthenticationResult) Preconditions.checkNotNull( + root.get(PlannerContext.DATA_CTX_AUTHENTICATION_RESULT), + "authenticationResult in dataContext" + ); + SystemSchema.checkStateReadAccessForServers(authenticationResult, authorizerMapper); + + final Set<String> serverFilter = SystemSchemaFilters.extractColumnValues(filters, SERVER_INDEX); + InvalidInput.conditionalException( + serverFilter != null, + "sys.stack_trace requires a filter on the server column using '=' or 'IN'" + ); + final int maxStackTraceFrameDepth = getMaxStackTraceFrameDepth( + root.get(StackTraceCollector.MAX_STACK_TRACE_FRAME_DEPTH_KEY) + ); + + final Iterator<DiscoveryDruidNode> druidServers = SystemSchema.getDruidServers(druidNodeDiscoveryProvider); + final Map<String, ServerStackTraceTarget> serverToTargetMap = new HashMap<>(); + druidServers.forEachRemaining(discoveryDruidNode -> { + final DruidNode druidNode = discoveryDruidNode.getDruidNode(); + final String server = druidNode.getHostAndPortToUse(); + if (!serverFilter.contains(server)) { + return; + } + + final String nodeRole = discoveryDruidNode.getNodeRole().getJsonName(); + final ServerStackTraceTarget target = serverToTargetMap.get(server); + if (target == null) { + serverToTargetMap.put( + server, + new ServerStackTraceTarget( + server, + druidNode.getServiceName(), + new ArrayList<>(Collections.singletonList(nodeRole)), + druidNode + ) + ); + } else { + target.addNodeRole(nodeRole); + } + }); + + final List<Object[]> rows = new ArrayList<>(); + for (final ServerStackTraceTarget target : serverToTargetMap.values()) { + rows.addAll(target.buildRows(this, projects, maxStackTraceFrameDepth)); + } + return Linq4j.asEnumerable(rows); + } + + static int getMaxStackTraceFrameDepth(@Nullable final Object value) + { + return StackTraceCollector.validateMaxStackTraceFrameDepth( + QueryContexts.getAsLong( + StackTraceCollector.MAX_STACK_TRACE_FRAME_DEPTH_KEY, + value, + StackTraceCollector.DEFAULT_MAX_STACK_TRACE_FRAME_DEPTH + ) + ); + } + + private static Object[] projectRow(final Object[] row, @Nullable final int[] projects) + { + if (projects == null) { + return row; + } + final Object[] projectedRow = new Object[projects.length]; + for (int i = 0; i < projects.length; i++) { + projectedRow[i] = row[projects[i]]; + } + return projectedRow; + } + + private StackTraceResult getStackTrace( + final DruidNode druidNode, + final int maxStackTraceFrameDepth + ) + { + final String url = druidNode.getUriToUse().resolve( + StringUtils.format( + "/status/stack?%s=%d", + StackTraceCollector.MAX_STACK_TRACE_FRAME_DEPTH_KEY, + maxStackTraceFrameDepth + ) + ).toString(); + try { + final Request request = new Request(HttpMethod.GET, URI.create(url).toURL()); + final StringFullResponseHolder response = httpClient + .go(request, new StringFullResponseHandler(StandardCharsets.UTF_8)) Review Comment: [P2] Cancel the HTTP request when SQL execution is interrupted Calling `.get()` directly discards the request future. If a `sys.stack_trace` query is canceled while a node is returning a large or deep snapshot, the catch block interrupts the caller but leaves the underlying HTTP request running and buffering until completion or timeout. Retain the future and use `FutureUtils.get(future, true)`, or explicitly cancel it in the interruption path. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
