This is an automated email from the ASF dual-hosted git repository. FrankChen021 pushed a commit to branch codex/native-sys-servers in repository https://gitbox.apache.org/repos/asf/druid.git
commit 2e00ba4803890d1e88934f1cd6834f4a87ff8439 Author: Frank Chen <[email protected]> AuthorDate: Wed Sep 2 18:14:32 2026 +0800 feat(sql): route local system tables in process --- docs/querying/sql-metadata-tables.md | 2 +- .../system/handler/SystemTableNodeLocator.java | 6 +- .../system/handler/SystemTableQueryClient.java | 35 ++++++++++- .../system/table/ServersTableDescriptor.java | 4 +- .../server/system/table/SystemTableDescriptor.java | 2 +- .../system/table/SystemTableRoutingMode.java | 3 + .../system/ServersTableDataProviderTest.java | 6 +- .../system/handler/SystemTableQueryClientTest.java | 71 ++++++++++++++++++++++ 8 files changed, 120 insertions(+), 9 deletions(-) diff --git a/docs/querying/sql-metadata-tables.md b/docs/querying/sql-metadata-tables.md index 0ffa44fc47c..0b8a56c4030 100644 --- a/docs/querying/sql-metadata-tables.md +++ b/docs/querying/sql-metadata-tables.md @@ -170,7 +170,7 @@ execution: |-----|--------------| |[`sys.tasks`](#tasks-table)|The Overlord that owns task state. Supported filters are pushed into task storage when the configured task storage implementation supports filter pushdown.| |[`sys.server_properties`](#server_properties-table)|The Druid server processes discovered in the cluster. Filters on `server` and `service_name` can avoid reading properties from nodes that don't match.| -|[`sys.servers`](#servers-table)|The current Coordinator leader's discovered-cluster view of Druid servers.| +|[`sys.servers`](#servers-table)|The Broker's local discovered-cluster view of Druid servers.| After Druid retrieves the system-table rows, the native engine applies the remaining filters, expressions, aggregations, sorting, and result processing. A system table that doesn't advertise native query support continues to diff --git a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNodeLocator.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNodeLocator.java index 1706860c140..ad6fcb641d2 100644 --- a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNodeLocator.java +++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNodeLocator.java @@ -68,7 +68,11 @@ public class SystemTableNodeLocator List<SystemTableNode> locate(final SystemTableDescriptor descriptor, final Query<?> query) { - if (descriptor.getRoutingMode() == SystemTableRoutingMode.ALL_NODES) { + final SystemTableRoutingMode routingMode = descriptor.getRoutingMode(); + if (routingMode == SystemTableRoutingMode.LOCAL) { + throw new ISE("Local system-table routing does not select a remote node"); + } + if (routingMode == SystemTableRoutingMode.ALL_NODES) { return discoverAllNodes(descriptor); } diff --git a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java index 7330e51553e..00613b1f7c8 100644 --- a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java +++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java @@ -282,13 +282,21 @@ public class SystemTableQueryClient implements DataSourceQueryHandler ) { final List<QueryRunner<ScanResultValue>> nodeRunners = new ArrayList<>(); + final SystemTableRoutingMode routingMode = descriptor.getRoutingMode(); + if (routingMode == SystemTableRoutingMode.LOCAL) { + nodeRunners.add( + (queryPlus, responseContext) -> runLocalNodeQuery(nodeQuery, queryPlus, responseContext) + ); + return nodeRunners; + } + for (final SystemTableNode node : nodeLocator.locate(descriptor, nodeQuery)) { nodeRunners.add( (queryPlus, responseContext) -> recoverNodeFailure( () -> runNodeQuery(nodeQuery, node, nodeSequencesCloser, queryPlus, responseContext), node, descriptor, - descriptor.getRoutingMode() == SystemTableRoutingMode.LEADER_ONLY + routingMode == SystemTableRoutingMode.LEADER_ONLY ? () -> { final SystemTableNode currentLeader = Iterables.getOnlyElement( nodeLocator.locate(descriptor, nodeQuery) @@ -305,6 +313,31 @@ public class SystemTableQueryClient implements DataSourceQueryHandler return nodeRunners; } + private Sequence<ScanResultValue> runLocalNodeQuery( + final ScanQuery nodeQuery, + final QueryPlus<ScanResultValue> queryPlus, + final ResponseContext responseContext + ) + { + final String nodeResourceId = UUID.randomUUID().toString(); + final String nodeQueryId = SystemTableDataSource.NODE_QUERY_ID_PREFIX + UUID.randomUUID(); + final ScanQuery subNativeQuery = nodeQuery.withOverriddenContext( + Map.of( + BaseQuery.QUERY_ID, + nodeQueryId, + QueryContexts.QUERY_RESOURCE_ID, + nodeResourceId + ) + ); + // Invoke the raw local handler rather than the Broker handler so local routing cannot recurse into node routing. + final QueryRunner<ScanResultValue> nodeRunner = localQueryHandler.createRunner( + subNativeQuery, + escalatedAuthenticationResult, + true + ); + return nodeRunner.run(queryPlus.withQuery(subNativeQuery), responseContext); + } + private Sequence<ScanResultValue> runNodeQuery( final ScanQuery nodeQuery, final SystemTableNode node, diff --git a/server/src/main/java/org/apache/druid/server/system/table/ServersTableDescriptor.java b/server/src/main/java/org/apache/druid/server/system/table/ServersTableDescriptor.java index e198dbeed23..dd7329a4b03 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/ServersTableDescriptor.java +++ b/server/src/main/java/org/apache/druid/server/system/table/ServersTableDescriptor.java @@ -56,7 +56,7 @@ public class ServersTableDescriptor implements SystemTableDescriptor .add("total_memory", ColumnType.LONG) .build(); - private static final Set<NodeRole> NODE_ROLES = Set.of(NodeRole.COORDINATOR); + private static final Set<NodeRole> NODE_ROLES = Set.of(NodeRole.BROKER); private static final SystemTableRowAuthorizer ROW_AUTHORIZER = (rows, authenticationResult, authorizerMapper) -> { final AuthorizationResult authorizationResult = AuthorizationUtils.authorizeAllResourceActions( authenticationResult, @@ -84,7 +84,7 @@ public class ServersTableDescriptor implements SystemTableDescriptor @Override public SystemTableRoutingMode getRoutingMode() { - return SystemTableRoutingMode.LEADER_ONLY; + return SystemTableRoutingMode.LOCAL; } @Override diff --git a/server/src/main/java/org/apache/druid/server/system/table/SystemTableDescriptor.java b/server/src/main/java/org/apache/druid/server/system/table/SystemTableDescriptor.java index f9442c51327..a7c85c7fd49 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/SystemTableDescriptor.java +++ b/server/src/main/java/org/apache/druid/server/system/table/SystemTableDescriptor.java @@ -33,7 +33,7 @@ public interface SystemTableDescriptor /** * Returns the node roles capable of serving this table. {@link #getRoutingMode()} determines whether every node or - * only the leader for the role receives a query. + * only the leader for the role receives a query, or whether the Broker runs the provider locally. */ Set<NodeRole> getNodeRoles(); diff --git a/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java b/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java index fa6c1c69ba6..5add89cff2b 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java +++ b/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java @@ -22,6 +22,9 @@ package org.apache.druid.server.system.table; /** Defines how the Broker selects nodes that contribute rows to a system table. */ public enum SystemTableRoutingMode { + /** The Broker executes the table provider locally without selecting a remote node. */ + LOCAL, + /** Every discovered node for the descriptor's roles contributes an independent set of rows. */ ALL_NODES, diff --git a/server/src/test/java/org/apache/druid/server/system/ServersTableDataProviderTest.java b/server/src/test/java/org/apache/druid/server/system/ServersTableDataProviderTest.java index 1fa2e4df385..1fe59067da7 100644 --- a/server/src/test/java/org/apache/druid/server/system/ServersTableDataProviderTest.java +++ b/server/src/test/java/org/apache/druid/server/system/ServersTableDataProviderTest.java @@ -159,12 +159,12 @@ public class ServersTableDataProviderTest } @Test - public void testDescriptorRoutesToCoordinatorLeader() + public void testDescriptorRunsLocallyOnBroker() { final ServersTableDescriptor descriptor = new ServersTableDescriptor(); - Assertions.assertEquals(Set.of(NodeRole.COORDINATOR), descriptor.getNodeRoles()); - Assertions.assertEquals(SystemTableRoutingMode.LEADER_ONLY, descriptor.getRoutingMode()); + Assertions.assertEquals(Set.of(NodeRole.BROKER), descriptor.getNodeRoles()); + Assertions.assertEquals(SystemTableRoutingMode.LOCAL, descriptor.getRoutingMode()); Assertions.assertEquals(ColumnType.NESTED_DATA, descriptor.getRowSignature().getColumnType(13).orElseThrow()); } diff --git a/server/src/test/java/org/apache/druid/server/system/handler/SystemTableQueryClientTest.java b/server/src/test/java/org/apache/druid/server/system/handler/SystemTableQueryClientTest.java index d5931fcbb5e..9fa21ab59a8 100644 --- a/server/src/test/java/org/apache/druid/server/system/handler/SystemTableQueryClientTest.java +++ b/server/src/test/java/org/apache/druid/server/system/handler/SystemTableQueryClientTest.java @@ -83,6 +83,7 @@ import org.apache.druid.server.security.NoopEscalator; import org.apache.druid.server.system.SystemTableNotLeaderException; import org.apache.druid.server.system.table.ServerPropertiesTableDescriptor; import org.apache.druid.server.system.table.SystemTableDescriptor; +import org.apache.druid.server.system.table.SystemTableRoutingMode; import org.apache.druid.server.system.table.TaskTableDescriptor; import org.joda.time.Duration; import org.junit.jupiter.api.Assertions; @@ -340,6 +341,59 @@ public class SystemTableQueryClientTest ); } + /** A LOCAL system table executes its provider in-process without node discovery or an HTTP client. */ + @Test + public void testLocalRoutingExecutesLocalRunnerWithoutHttp() + { + final SystemTableDescriptor descriptor = new TestSystemTableDescriptor(SystemTableRoutingMode.LOCAL); + final SystemTableNodeLocator nodeLocator = Mockito.mock(SystemTableNodeLocator.class); + final DirectDruidClientFactory directClientFactory = Mockito.mock(DirectDruidClientFactory.class); + final SystemTableQueryHandler localQueryHandler = Mockito.mock(SystemTableQueryHandler.class); + final AuthenticationResult escalatedAuthenticationResult = + new AuthenticationResult("system", "allow", "system", null); + final Escalator escalator = Mockito.mock(Escalator.class); + Mockito.when(escalator.createEscalatedAuthenticationResult()).thenReturn(escalatedAuthenticationResult); + Mockito.doAnswer( + ignored -> (QueryRunner<ScanResultValue>) (queryPlus, responseContext) -> + Sequences.simple(List.of(scanResult(new Object[]{"local"}))) + ).when(localQueryHandler).createRunner( + ArgumentMatchers.any(), + Mockito.same(escalatedAuthenticationResult), + Mockito.eq(true) + ); + + final QuerySegmentWalker querySegmentWalker = Mockito.mock(QuerySegmentWalker.class); + Mockito.when(querySegmentWalker.getQueryRunnerForIntervals(ArgumentMatchers.any(), ArgumentMatchers.any())) + .thenAnswer(ignored -> passthroughRunner(descriptor.getRowSignature())); + final SystemTableQueryClient client = new SystemTableQueryClient( + nodeLocator, + directClientFactory, + Mockito.mock(QueryScheduler.class), + querySegmentWalker, + Map.of(descriptor.getTableName(), descriptor), + new AuthorizerMapper(Map.of("allow", new AllowAllAuthorizer(null))), + localQueryHandler, + escalator, + nonMatchingSelfNode() + ); + final ScanQuery query = query(descriptor, Collections.emptyMap()); + + final List<ScanResultValue> results = client.createRunner(query, AUTHENTICATION_RESULT, false) + .run(QueryPlus.wrap(query), ResponseContext.createEmpty()) + .toList(); + + Assertions.assertEquals(1, results.size()); + final List<?> events = (List<?>) results.get(0).getEvents(); + Assertions.assertEquals(1, events.size()); + Assertions.assertArrayEquals(new Object[]{"local"}, (Object[]) events.get(0)); + Mockito.verifyNoInteractions(nodeLocator, directClientFactory); + Mockito.verify(localQueryHandler).createRunner( + ArgumentMatchers.any(), + Mockito.same(escalatedAuthenticationResult), + Mockito.eq(true) + ); + } + /** Delayed node responses verify concurrent fanout through the real {@link DirectDruidClient} transport path. */ @Test public void testDirectClientsStartRequestsBeforeWaitingForDelayedResponses() throws Exception @@ -1382,6 +1436,17 @@ public class SystemTableQueryClientTest private static class TestSystemTableDescriptor implements SystemTableDescriptor { private static final RowSignature ROW_SIGNATURE = RowSignature.builder().add("value", null).build(); + private final SystemTableRoutingMode routingMode; + + private TestSystemTableDescriptor() + { + this(SystemTableRoutingMode.ALL_NODES); + } + + private TestSystemTableDescriptor(final SystemTableRoutingMode routingMode) + { + this.routingMode = routingMode; + } @Override public String getTableName() @@ -1395,6 +1460,12 @@ public class SystemTableQueryClientTest return Set.of(NodeRole.BROKER); } + @Override + public SystemTableRoutingMode getRoutingMode() + { + return routingMode; + } + @Override public RowSignature getRowSignature() { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
