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]

Reply via email to