diff --git a/docs/querying/sql-metadata-tables.md b/docs/querying/sql-metadata-tables.md
index 055305d81e50..8b47335617fe 100644
--- a/docs/querying/sql-metadata-tables.md
+++ b/docs/querying/sql-metadata-tables.md
@@ -135,10 +135,57 @@ WHERE "IS_AGGREGATOR" = 'YES'
The "sys" schema provides visibility into Druid segments, servers and tasks.
:::info
- Note: "sys" tables do not currently support Druid-specific functions like `TIME_PARSE` and
- `APPROX_QUANTILE_DS`. Only standard SQL functions can be used.
+ By default, "sys" tables use the SQL-layer execution path and support only standard SQL functions. You can enable
+ [native query execution](#native-query-execution) for supported system tables, which enables expressions and
+ aggregations that the native SQL engine can translate.
:::
+### Native query execution
+
+The native SQL engine can plan supported system tables as native datasources. To enable this behavior for a query, set
+`useNativeQueryForSystemTables` to `true` in the SQL query context:
+
+```json
+{
+ "query": "SELECT COUNT(DISTINCT server) FROM sys.server_properties",
+ "context": {
+ "useNativeQueryForSystemTables": true
+ }
+}
+```
+
+For clients that support setting query context parameters with SQL statements, you can enable native system-table
+execution with `SET`:
+
+```sql
+SET useNativeQueryForSystemTables = 'true';
+SELECT COUNT(DISTINCT server) FROM sys.server_properties;
+```
+
+Native system-table execution is available when the resolved SQL engine is `native`. You don't need to explicitly set
+the `engine` context parameter when `native` is already the default engine. The following tables support native query
+execution:
+
+|Table|Source of rows|
+|-----|--------------|
+|[`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.|
+
+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
+use its existing SQL-layer execution path, even when `useNativeQueryForSystemTables` is `true`.
+
+The parameter defaults to `false`. During a rolling upgrade, leave it disabled until the Broker and all nodes
+that serve the native system tables have been upgraded. After the upgrade, you can enable it by query or set
+`druid.query.default.context.useNativeQueryForSystemTables=true` on Brokers as the cluster-wide default. For more
+information, see [SQL query context](sql-query-context.md).
+
+Native system-table queries sent to the Router use distributed Broker execution by default. To execute a native
+system-table Scan against only the contacted node, set the HTTP header
+`X-Druid-Native-Query-Route: local`. Local execution uses the authenticated request identity and applies the table's
+authorization rules. The Broker uses the same header for remote node fan-out requests. If the Broker itself is
+one of the selected nodes, it executes that node Scan in-process without an HTTP request. The header
+controls routing and doesn't grant additional permissions.
+
### SEGMENTS table
Segments table provides details on all Druid segments, whether they are published yet or not.
diff --git a/docs/querying/sql-query-context.md b/docs/querying/sql-query-context.md
index 88cafe1ced71..b5f813d17716 100644
--- a/docs/querying/sql-query-context.md
+++ b/docs/querying/sql-query-context.md
@@ -45,6 +45,7 @@ The table below lists the query context parameters you can use with Druid SQL.
|`useLexicographicTopN`|If `true`, Druid can use [TopN queries](topnquery.md) with lexicographic dimension ordering. If `false`, Druid uses [GroupBy queries](groupbyquery.md) instead for lexicographic ordering. When both `useLexicographicTopN` and `useApproximateTopN` are `false`, TopN queries are never used.|`false`|
|`enableTimeBoundaryPlanning`|If `true`, Druid converts SQL queries to [time boundary queries](timeboundaryquery.md) wherever possible. Time boundary queries are very efficient for min-max calculation on the `__time` column in a datasource. |`false`|
|`useNativeQueryExplain`|If `true`, `EXPLAIN PLAN FOR` returns the explain plan as a JSON representation of equivalent native query, else it returns the original version of explain plan generated by Calcite.
This property is provided for backwards compatibility. We don't recommend setting this parameter unless your application depends on the older behavior.|`true`|
+|`useNativeQueryForSystemTables`|If `true` and the resolved SQL engine is `native`, Druid uses native query planning for [system tables](sql-metadata-tables.md#native-query-execution) that advertise native query support. System tables without native query support continue to use their existing execution path.
During a rolling upgrade, leave this parameter disabled until the Broker and the components that serve the native system tables have been upgraded. You can set a cluster-wide default on Brokers with `druid.query.default.context.useNativeQueryForSystemTables`.|`false`|
|`sqlFinalizeOuterSketches`|If `false` (default behavior in Druid 25.0.0 and later), `DS_HLL`, `DS_THETA`, and `DS_QUANTILES_SKETCH` return sketches in query results. If `true` (default behavior in Druid 24.0.1 and earlier), Druid finalizes sketches from these functions when they appear in query results.
This property is provided for backwards compatibility with behavior in Druid 24.0.1 and earlier. We don't recommend setting this parameter unless your application uses Druid 24.0.1 or earlier. Instead, use a function that doesn't return a sketch, such as `APPROX_COUNT_DISTINCT_DS_HLL`, `APPROX_COUNT_DISTINCT_DS_THETA`, `APPROX_QUANTILE_DS`, `DS_THETA_ESTIMATE`, or `DS_GET_QUANTILE`.|`false`|
|`sqlUseBoundAndSelectors`|If `false` (default behavior in Druid 27.0.0 and later), the SQL planner uses [equality](./filters.md#equality-filter), [null](./filters.md#null-filter), and [range](./filters.md#range-filter) filters instead of [selector](./filters.md#selector-filter) and [bounds](./filters.md#bound-filter). For filtering `ARRAY` typed values, `sqlUseBoundAndSelectors` must be `false`. | `false`.|
|`sqlUseExtractionFns`|If false, the SQL planner avoids using [`extractionFn`](dimensionspecs.md#extraction-functions) in favor of using other constructs such as [virtual columns](virtual-columns.md). This parameter is provided for compatibility with prior behavior, and may be removed in a future release.|false|
diff --git a/docs/querying/sql-translation.md b/docs/querying/sql-translation.md
index 8c2e323b6b48..3c3384148332 100644
--- a/docs/querying/sql-translation.md
+++ b/docs/querying/sql-translation.md
@@ -817,14 +817,20 @@ This query context parameter is a temporary solution to avoid the known issue.
Druid does not support all SQL features. In particular, the following features are not supported.
-- JOIN between native datasources (table, lookup, subquery) and [system tables](sql-metadata-tables.md).
+- JOIN between native datasources (table, lookup, subquery) and [system tables](sql-metadata-tables.md) that use the
+ traditional SQL-layer execution path. System tables that support native execution can participate in these joins
+ when `useNativeQueryForSystemTables` is enabled.
- JOIN conditions that are not an equality between expressions from the left- and right-hand sides.
- JOIN conditions containing a constant value inside the condition.
- JOIN conditions on a column which contains a multi-value dimension.
-- ORDER BY for a non-aggregating query, except for `ORDER BY __time` or `ORDER BY __time DESC`, which are supported.
- This restriction only applies to non-aggregating queries; you can ORDER BY any column in an aggregating query.
+- ORDER BY for a non-aggregating query, except for `ORDER BY __time` or `ORDER BY __time DESC`, and native-enabled
+ system-table queries that use the native window-query path. This restriction only applies to non-aggregating queries;
+ you can ORDER BY any column in an aggregating query.
- DDL and DML.
-- Using Druid-specific functions like `TIME_PARSE` and `APPROX_QUANTILE_DS` on [system tables](sql-metadata-tables.md).
+- Using Druid-specific functions like `TIME_PARSE` and `APPROX_QUANTILE_DS` on system tables that use the traditional
+ SQL-layer execution path. [System tables with native query support](sql-metadata-tables.md#native-query-execution) can
+ use expressions and aggregations that the native SQL engine can translate when `useNativeQueryForSystemTables` is
+ enabled.
Additionally, some Druid native query features are not supported by the SQL language. Some unsupported Druid features
include:
@@ -834,5 +840,3 @@ include:
- [Multi-value dimensions](sql-data-types.md#multi-value-strings) are only partially implemented in Druid SQL. There are known
inconsistencies between their behavior in SQL queries and in native queries due to how they are currently treated by
the SQL planner.
-
-
diff --git a/docs/querying/sql.md b/docs/querying/sql.md
index b74c36365025..f69b7affb109 100644
--- a/docs/querying/sql.md
+++ b/docs/querying/sql.md
@@ -79,11 +79,15 @@ datasources can be referenced as either `druid.dataSourceName` or simply `dataSo
- [Lookups](datasource.md#lookup) from the `lookup` schema, for example `lookup.countries`. Note that lookups can
also be queried using the [`LOOKUP` function](sql-scalar.md#string-functions).
- [Subqueries](datasource.md#query).
-- [Joins](datasource.md#join) between anything in this list, except between native datasources (table, lookup,
-query) and system tables. The join condition must be an equality between expressions from the left- and right-hand side
-of the join.
-- [Metadata tables](sql-metadata-tables.md) from the `INFORMATION_SCHEMA` or `sys` schemas. Unlike the other options for the
-FROM clause, metadata tables are not considered datasources. They exist only in the SQL layer.
+- [Joins](datasource.md#join) between anything in this list. The traditional SQL-layer system-table path does not
+ support joins between native datasources (table, lookup, query) and system tables. When
+ [`useNativeQueryForSystemTables`](sql-metadata-tables.md#native-query-execution) is enabled, system tables that support
+ native execution can participate in native datasource joins. The join condition must be an equality between
+ expressions from the left- and right-hand side of the join.
+- [Metadata tables](sql-metadata-tables.md) from the `INFORMATION_SCHEMA` or `sys` schemas. By default, metadata tables
+ exist only in the SQL layer and are not considered datasources. When
+ [`useNativeQueryForSystemTables`](sql-metadata-tables.md#native-query-execution) is enabled, supported `sys` tables are
+ represented as native system-table datasources.
For more information about table, lookup, query, and join datasources, refer to the [Datasources](datasource.md)
documentation.
@@ -295,7 +299,8 @@ grouping expressions or aggregated values. It can only be used together with GRO
The ORDER BY clause refers to columns that are present after execution of GROUP BY. It can be used to order the results
based on either grouping expressions or aggregated values. ORDER BY can refer to an expression or a select clause
ordinal position (like `ORDER BY 2` to order by the second selected column). For non-aggregation queries, ORDER BY
-can only order by the `__time` column. For aggregation queries, ORDER BY can order by any column.
+can normally only order by the `__time` column. Native-enabled system tables can also use the native window-query path
+to order by other columns. For aggregation queries, ORDER BY can order by any column.
## LIMIT
diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerPropertiesQueryTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerPropertiesQueryTest.java
new file mode 100644
index 000000000000..18334a0547fa
--- /dev/null
+++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerPropertiesQueryTest.java
@@ -0,0 +1,182 @@
+/*
+ * 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.testing.embedded.query;
+
+import org.apache.druid.query.QueryContexts;
+import org.apache.druid.sql.calcite.planner.PlannerContext;
+import org.apache.druid.sql.calcite.run.NativeSqlEngine;
+import org.apache.druid.testing.embedded.EmbeddedBroker;
+import org.apache.druid.testing.embedded.EmbeddedCoordinator;
+import org.apache.druid.testing.embedded.EmbeddedDruidCluster;
+import org.apache.druid.testing.embedded.EmbeddedHistorical;
+import org.apache.druid.testing.embedded.EmbeddedIndexer;
+import org.apache.druid.testing.embedded.EmbeddedOverlord;
+import org.apache.druid.testing.embedded.EmbeddedRouter;
+import org.apache.druid.testing.embedded.junit5.EmbeddedClusterTestBase;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.util.Map;
+
+public class NativeSysServerPropertiesQueryTest extends EmbeddedClusterTestBase
+{
+ private static final String SERVICE_NAME = "native/mvp/broker";
+ private static final String COORDINATOR_PROPERTY = "native.sys.server.properties.coordinator";
+ private static final String OVERLORD_PROPERTY = "native.sys.server.properties.overlord";
+ private static final String BROKER_PROPERTY = "native.sys.server.properties.broker";
+ private static final String HISTORICAL_PROPERTY = "native.sys.server.properties.historical";
+ private static final String INDEXER_PROPERTY = "native.sys.server.properties.indexer";
+ private static final String ROUTER_PROPERTY = "native.sys.server.properties.router";
+
+ private final EmbeddedCoordinator coordinator = new EmbeddedCoordinator()
+ .addProperty(COORDINATOR_PROPERTY, "enabled");
+
+ private final EmbeddedOverlord overlord = new EmbeddedOverlord()
+ .addProperty(OVERLORD_PROPERTY, "enabled");
+
+ private final EmbeddedBroker broker = new EmbeddedBroker()
+ .addProperty("druid.service", SERVICE_NAME)
+ .addProperty(BROKER_PROPERTY, "enabled");
+
+ private final EmbeddedHistorical historical = new EmbeddedHistorical()
+ .addProperty(HISTORICAL_PROPERTY, "enabled");
+
+ private final EmbeddedIndexer indexer = new EmbeddedIndexer()
+ .addProperty(INDEXER_PROPERTY, "enabled");
+
+ private final EmbeddedRouter router = new EmbeddedRouter()
+ .addProperty(ROUTER_PROPERTY, "enabled");
+
+ @Override
+ protected EmbeddedDruidCluster createCluster()
+ {
+ return EmbeddedDruidCluster.withEmbeddedDerbyAndZookeeper()
+ .useLatchableEmitter()
+ .addCommonProperty("druid.centralizedDatasourceSchema.enabled", "true")
+ .addServer(coordinator)
+ .addServer(overlord)
+ .addServer(broker)
+ .addServer(historical)
+ .addServer(indexer)
+ .addServer(router);
+ }
+
+ /**
+ * Verifies that a native query for {@code sys.server_properties} reaches every persistent node type in the
+ * embedded cluster. In particular, the Router row proves that the Broker's native request included
+ * {@code X-Druid-Native-Query-Route: local}: without that header the Router would forward the request back to the
+ * Broker instead of reading its local system-table provider.
+ */
+ @ParameterizedTest(name = "plannerStrategy = {0}")
+ @ValueSource(strings = {
+ QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_COUPLED,
+ QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_DECOUPLED
+ })
+ public void testServerPropertiesFansOutToAllNodes(final String plannerStrategy)
+ {
+ final String result = cluster.runSql(
+ "SELECT service_name, COUNT(*) "
+ + "FROM sys.server_properties "
+ + "WHERE property IN ('" + COORDINATOR_PROPERTY + "', '" + OVERLORD_PROPERTY + "', '" + BROKER_PROPERTY
+ + "', '" + HISTORICAL_PROPERTY + "', '" + INDEXER_PROPERTY + "', '" + ROUTER_PROPERTY + "') "
+ + "GROUP BY service_name ORDER BY service_name",
+ nativeQueryContext(plannerStrategy)
+ );
+
+ Assertions.assertEquals(
+ String.join(
+ "\n",
+ "druid/coordinator,1",
+ "druid/historical,1",
+ "druid/indexer,1",
+ "druid/overlord,1",
+ "druid/router,1",
+ SERVICE_NAME + ",1"
+ ),
+ result
+ );
+ }
+
+ @ParameterizedTest(name = "plannerStrategy = {0}")
+ @ValueSource(strings = {
+ QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_COUPLED,
+ QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_DECOUPLED
+ })
+ public void testNativeAggregationsSupportDistinctCount(final String plannerStrategy)
+ {
+ final String result = cluster.runSql(
+ "SELECT COUNT(*), COUNT(DISTINCT service_name), COUNT(DISTINCT server), "
+ + "COUNT(DISTINCT property), SUM(1) "
+ + "FROM sys.server_properties "
+ + "WHERE property IN ('" + COORDINATOR_PROPERTY + "', '" + OVERLORD_PROPERTY + "', '" + BROKER_PROPERTY
+ + "', '" + HISTORICAL_PROPERTY + "', '" + INDEXER_PROPERTY + "', '" + ROUTER_PROPERTY + "')",
+ nativeQueryContext(plannerStrategy)
+ );
+
+ Assertions.assertEquals("6,6,6,6,6", result);
+ }
+
+ /**
+ * A filtered window-function query is planned as a {@code WindowOperatorQuery} over a system-table Scan, whose
+ * filter is pushed into the node scans.
+ */
+ @ParameterizedTest(name = "plannerStrategy = {0}")
+ @ValueSource(strings = {
+ QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_COUPLED,
+ QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_DECOUPLED
+ })
+ public void testNativeWindowFunction(final String plannerStrategy)
+ {
+ final String result = cluster.runSql(
+ "SELECT service_name, ROW_NUMBER() OVER (ORDER BY service_name) "
+ + "FROM sys.server_properties "
+ + "WHERE property IN ('" + COORDINATOR_PROPERTY + "', '" + OVERLORD_PROPERTY + "', '" + BROKER_PROPERTY
+ + "', '" + HISTORICAL_PROPERTY + "', '" + INDEXER_PROPERTY + "', '" + ROUTER_PROPERTY + "') "
+ + "ORDER BY service_name",
+ nativeQueryContext(plannerStrategy)
+ );
+
+ Assertions.assertEquals(
+ String.join(
+ "\n",
+ "druid/coordinator,1",
+ "druid/historical,2",
+ "druid/indexer,3",
+ "druid/overlord,4",
+ "druid/router,5",
+ SERVICE_NAME + ",6"
+ ),
+ result
+ );
+ }
+
+ private static Map nativeQueryContext(final String plannerStrategy)
+ {
+ return Map.of(
+ QueryContexts.ENGINE,
+ NativeSqlEngine.NAME,
+ PlannerContext.CTX_USE_NATIVE_QUERY_FOR_SYSTEM_TABLES,
+ true,
+ QueryContexts.CTX_NATIVE_QUERY_SQL_PLANNING_MODE,
+ plannerStrategy
+ );
+ }
+}
diff --git a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfo.java b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfo.java
index 41e1ef8b7c9a..de7b67b0d5d1 100644
--- a/indexing-service/src/main/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfo.java
+++ b/indexing-service/src/main/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfo.java
@@ -24,6 +24,7 @@
import com.google.inject.Inject;
import org.apache.druid.indexing.overlord.DruidOverlord;
import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.query.SystemTableDataSource;
import org.apache.druid.server.http.RedirectInfo;
import java.net.URL;
@@ -35,7 +36,9 @@ public class OverlordRedirectInfo implements RedirectInfo
{
private static final Set LOCAL_PATHS = ImmutableSet.of(
"/druid/indexer/v1/leader",
- "/druid/indexer/v1/isLeader"
+ "/druid/indexer/v1/isLeader",
+ "/druid/v2",
+ "/druid/v2/"
);
private final DruidOverlord overlord;
@@ -49,7 +52,10 @@ public OverlordRedirectInfo(DruidOverlord overlord)
@Override
public boolean doLocal(String requestURI)
{
- return (requestURI != null && LOCAL_PATHS.contains(requestURI)) || overlord.isLeader();
+ return (requestURI != null
+ && (LOCAL_PATHS.contains(requestURI)
+ || requestURI.startsWith("/druid/v2/" + SystemTableDataSource.NODE_QUERY_ID_PREFIX)))
+ || overlord.isLeader();
}
@Override
diff --git a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfoTest.java b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfoTest.java
index 7dde7dab0232..dbcf81bb7e61 100644
--- a/indexing-service/src/test/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfoTest.java
+++ b/indexing-service/src/test/java/org/apache/druid/indexing/overlord/http/OverlordRedirectInfoTest.java
@@ -50,6 +50,7 @@ public void testDoLocalWhenLeading()
Assertions.assertTrue(redirectInfo.doLocal(null));
Assertions.assertTrue(redirectInfo.doLocal("/druid/indexer/v1/leader"));
Assertions.assertTrue(redirectInfo.doLocal("/druid/indexer/v1/isLeader"));
+ Assertions.assertTrue(redirectInfo.doLocal("/druid/v2"));
Assertions.assertTrue(redirectInfo.doLocal("/druid/indexer/v1/other/path"));
EasyMock.verify(overlord);
}
@@ -62,6 +63,7 @@ public void testDoLocalWhenNotLeading()
Assertions.assertFalse(redirectInfo.doLocal(null));
Assertions.assertTrue(redirectInfo.doLocal("/druid/indexer/v1/leader"));
Assertions.assertTrue(redirectInfo.doLocal("/druid/indexer/v1/isLeader"));
+ Assertions.assertTrue(redirectInfo.doLocal("/druid/v2"));
Assertions.assertFalse(redirectInfo.doLocal("/druid/indexer/v1/other/path"));
EasyMock.verify(overlord);
}
diff --git a/multi-stage-query/src/test/java/org/apache/druid/msq/dart/controller/http/DartSqlResourceTest.java b/multi-stage-query/src/test/java/org/apache/druid/msq/dart/controller/http/DartSqlResourceTest.java
index 58b4265430b3..a42cc4534cd8 100644
--- a/multi-stage-query/src/test/java/org/apache/druid/msq/dart/controller/http/DartSqlResourceTest.java
+++ b/multi-stage-query/src/test/java/org/apache/druid/msq/dart/controller/http/DartSqlResourceTest.java
@@ -562,6 +562,33 @@ public void test_doPost_informationSchema()
);
}
+ @Test
+ public void test_doPost_nativeCapableSystemTableUsesBindablePlan()
+ {
+ final MockAsyncContext asyncContext = new MockAsyncContext();
+ final MockHttpServletResponse asyncResponse = new MockHttpServletResponse();
+ asyncContext.response = asyncResponse;
+
+ Mockito.when(httpServletRequest.getAttribute(AuthConfig.DRUID_AUTHENTICATION_RESULT))
+ .thenReturn(makeAuthenticationResult(CalciteTests.TEST_SUPERUSER_NAME));
+ Mockito.when(httpServletRequest.startAsync())
+ .thenReturn(asyncContext);
+
+ final SqlQuery sqlQuery = new SqlQuery(
+ "SELECT COUNT(*) FROM sys.tasks",
+ ResultFormat.ARRAY,
+ false,
+ false,
+ false,
+ Map.of(QueryContexts.ENGINE, DartSqlEngine.NAME),
+ Collections.emptyList()
+ );
+
+ Assertions.assertNull(sqlResource.doPost(sqlQuery, httpServletRequest));
+ Assertions.assertEquals(Response.Status.OK.getStatusCode(), asyncResponse.getStatus());
+ Assertions.assertEquals("[[4]]\n", StringUtils.fromUtf8(asyncResponse.baos.toByteArray()));
+ }
+
@Test
public void test_doPost_sysTableJoinedToDatasource()
{
diff --git a/processing/src/main/java/org/apache/druid/query/DataSource.java b/processing/src/main/java/org/apache/druid/query/DataSource.java
index 9b0d65f2d83a..a89ec5a43a07 100644
--- a/processing/src/main/java/org/apache/druid/query/DataSource.java
+++ b/processing/src/main/java/org/apache/druid/query/DataSource.java
@@ -49,7 +49,8 @@
@JsonSubTypes.Type(value = GlobalTableDataSource.class, name = "globalTable"),
@JsonSubTypes.Type(value = UnnestDataSource.class, name = "unnest"),
@JsonSubTypes.Type(value = FilteredDataSource.class, name = "filter"),
- @JsonSubTypes.Type(value = RestrictedDataSource.class, name = "restrict")
+ @JsonSubTypes.Type(value = RestrictedDataSource.class, name = "restrict"),
+ @JsonSubTypes.Type(value = SystemTableDataSource.class, name = "systemTable")
})
public interface DataSource extends Cacheable
{
diff --git a/processing/src/main/java/org/apache/druid/query/SystemTableDataSource.java b/processing/src/main/java/org/apache/druid/query/SystemTableDataSource.java
new file mode 100644
index 000000000000..2d2886619669
--- /dev/null
+++ b/processing/src/main/java/org/apache/druid/query/SystemTableDataSource.java
@@ -0,0 +1,110 @@
+/*
+ * 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.query;
+
+import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonProperty;
+import com.google.common.base.Preconditions;
+
+import java.util.Collections;
+import java.util.Objects;
+import java.util.Set;
+
+/**
+ * Identifies a Druid system table whose rows are supplied by a service rather than by segments.
+ */
+public class SystemTableDataSource extends LeafDataSource
+{
+ public static final String NODE_QUERY_ID_PREFIX = "native-system-node-";
+
+ private final String table;
+
+ @JsonCreator
+ public SystemTableDataSource(@JsonProperty("table") final String table)
+ {
+ this.table = Preconditions.checkNotNull(table, "table");
+ }
+
+ @JsonProperty
+ public String getTable()
+ {
+ return table;
+ }
+
+ @Override
+ public Set getTableNames()
+ {
+ // QueryScheduler uses table names as cancellation authorization resources. Namespace the synthetic datasource
+ // name so a system table cannot collide with a regular Druid datasource that has the same unqualified name.
+ return Collections.singleton("sys." + table);
+ }
+
+ @Override
+ public boolean isCacheable(boolean isBroker)
+ {
+ return false;
+ }
+
+ @Override
+ public boolean isGlobal()
+ {
+ return false;
+ }
+
+ @Override
+ public boolean isProcessable()
+ {
+ // The owning service must first resolve this datasource to an InlineDataSource.
+ return false;
+ }
+
+ @Override
+ public byte[] getCacheKey()
+ {
+ return null;
+ }
+
+ @Override
+ public boolean equals(Object o)
+ {
+ if (this == o) {
+ return true;
+ }
+ if (o == null || getClass() != o.getClass()) {
+ return false;
+ }
+ final SystemTableDataSource that = (SystemTableDataSource) o;
+ return Objects.equals(table, that.table);
+ }
+
+ @Override
+ public int hashCode()
+ {
+ return Objects.hash(table);
+ }
+
+ @Override
+ public String toString()
+ {
+ return "SystemTableDataSource{" +
+ "table='" + table + '\'' +
+ '}';
+ }
+}
diff --git a/processing/src/main/java/org/apache/druid/query/rowsandcols/LazilyDecoratedRowsAndColumns.java b/processing/src/main/java/org/apache/druid/query/rowsandcols/LazilyDecoratedRowsAndColumns.java
index 270ab566814b..8fb513643cde 100644
--- a/processing/src/main/java/org/apache/druid/query/rowsandcols/LazilyDecoratedRowsAndColumns.java
+++ b/processing/src/main/java/org/apache/druid/query/rowsandcols/LazilyDecoratedRowsAndColumns.java
@@ -32,6 +32,7 @@
import org.apache.druid.java.util.common.ISE;
import org.apache.druid.java.util.common.Pair;
import org.apache.druid.java.util.common.UOE;
+import org.apache.druid.query.ResourceLimitExceededException;
import org.apache.druid.query.filter.Filter;
import org.apache.druid.query.filter.ValueMatcher;
import org.apache.druid.query.operator.ColumnWithDirection;
@@ -281,7 +282,9 @@ private Pair materializeCursorFactory(CursorFactory cursor
cursor.advance();
}
for (; !cursor.isDoneOrInterrupted() && remainingRowsToFetch > 0; remainingRowsToFetch--) {
- writer.addSelection();
+ if (!writer.addSelection()) {
+ throw materializationLimitExceeded();
+ }
cursor.advance();
}
@@ -395,10 +398,20 @@ private Pair naiveMaterialize(RowsAndColumns rac)
continue;
}
remainingRowsToFetch--;
- frameWriter.addSelection();
+ if (!frameWriter.addSelection()) {
+ throw materializationLimitExceeded();
+ }
}
return Pair.of(frameWriter.toByteArray(), sigBob.build());
}
}
+
+ private ResourceLimitExceededException materializationLimitExceeded()
+ {
+ return ResourceLimitExceededException.withMessage(
+ "RowsAndColumns materialization exceeded the configured frame capacity of [%,d] bytes",
+ allocatorCapacity
+ );
+ }
}
diff --git a/processing/src/main/java/org/apache/druid/segment/FilteredSegment.java b/processing/src/main/java/org/apache/druid/segment/FilteredSegment.java
index 909d80b38aaa..1a56bc48a3c7 100644
--- a/processing/src/main/java/org/apache/druid/segment/FilteredSegment.java
+++ b/processing/src/main/java/org/apache/druid/segment/FilteredSegment.java
@@ -20,6 +20,7 @@
package org.apache.druid.segment;
import org.apache.druid.query.filter.DimFilter;
+import org.apache.druid.query.rowsandcols.CursorFactoryRowsAndColumns;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
@@ -44,6 +45,8 @@ public T as(@Nonnull Class clazz)
{
if (CursorFactory.class.equals(clazz)) {
return (T) new FilteredCursorFactory(delegate.as(CursorFactory.class), filter);
+ } else if (CloseableShapeshifter.class.equals(clazz)) {
+ return (T) new CursorFactoryRowsAndColumns(as(CursorFactory.class));
} else if (TopNOptimizationInspector.class.equals(clazz)) {
return (T) new SimpleTopNOptimizationInspector(filter == null);
}
diff --git a/processing/src/test/java/org/apache/druid/query/rowsandcols/semantic/EvaluateRowsAndColumnsTest.java b/processing/src/test/java/org/apache/druid/query/rowsandcols/semantic/EvaluateRowsAndColumnsTest.java
index 810b1b83da21..eef947598c95 100644
--- a/processing/src/test/java/org/apache/druid/query/rowsandcols/semantic/EvaluateRowsAndColumnsTest.java
+++ b/processing/src/test/java/org/apache/druid/query/rowsandcols/semantic/EvaluateRowsAndColumnsTest.java
@@ -20,6 +20,7 @@
package org.apache.druid.query.rowsandcols.semantic;
import com.google.common.collect.Lists;
+import org.apache.druid.query.ResourceLimitExceededException;
import org.apache.druid.query.expression.TestExprMacroTable;
import org.apache.druid.query.operator.OffsetLimit;
import org.apache.druid.query.operator.window.RowsAndColumnsHelper;
@@ -118,4 +119,34 @@ public void testMaterializeColumns()
.expectColumn("array", expectedArr, ColumnType.STRING_ARRAY)
.validate(ras);
}
+
+ @Test
+ public void testMaterializationFailsInsteadOfReturningPartialRowsWhenFrameIsFull()
+ {
+ final Object[][] vals = new Object[][] {{0L, "x".repeat(64 * 1024)}};
+ final RowSignature signature = RowSignature.builder()
+ .add("__time", ColumnType.LONG)
+ .add("value", ColumnType.STRING)
+ .build();
+ final RowsAndColumns base = make(MapOfColumnsRowsAndColumns.fromRowObjects(vals, signature));
+
+ Assumptions.assumeTrue(base.as(CursorFactory.class) != null, "skipping: CursorFactory not supported");
+
+ final LazilyDecoratedRowsAndColumns rowsAndColumns = new LazilyDecoratedRowsAndColumns(
+ base,
+ null,
+ null,
+ null,
+ OffsetLimit.limit(1),
+ null,
+ null,
+ 1024L
+ );
+
+ final ResourceLimitExceededException exception = Assertions.assertThrows(
+ ResourceLimitExceededException.class,
+ rowsAndColumns::numRows
+ );
+ Assertions.assertTrue(exception.getMessage().contains("configured frame capacity"));
+ }
}
diff --git a/processing/src/test/java/org/apache/druid/segment/FilteredSegmentTest.java b/processing/src/test/java/org/apache/druid/segment/FilteredSegmentTest.java
new file mode 100644
index 000000000000..39cbda153464
--- /dev/null
+++ b/processing/src/test/java/org/apache/druid/segment/FilteredSegmentTest.java
@@ -0,0 +1,56 @@
+/*
+ * 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.segment;
+
+import org.apache.druid.query.InlineDataSource;
+import org.apache.druid.query.filter.SelectorDimFilter;
+import org.apache.druid.query.rowsandcols.RowsAndColumns;
+import org.apache.druid.segment.column.ColumnType;
+import org.apache.druid.segment.column.RowSignature;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.List;
+
+public class FilteredSegmentTest
+{
+ /** Window queries read segments as {@link RowsAndColumns}; a filtered segment provides them with its filter applied. */
+ @Test
+ public void testShapeshiftsToFilteredRowsAndColumns() throws Exception
+ {
+ final RowSignature signature = RowSignature.builder().add("dim", ColumnType.STRING).build();
+ final InlineDataSource inline = InlineDataSource.fromIterable(
+ List.of(new Object[]{"a"}, new Object[]{"b"}, new Object[]{"a"}),
+ signature
+ );
+ final Segment segment = new FilteredSegment(
+ new ArrayListSegment<>(new ArrayList<>(inline.getRowsAsList()), inline.rowAdapter(), signature),
+ new SelectorDimFilter("dim", "b", null)
+ );
+
+ try (final CloseableShapeshifter shapeshifter = segment.as(CloseableShapeshifter.class)) {
+ Assertions.assertNotNull(shapeshifter);
+ final RowsAndColumns rac = Assertions.assertInstanceOf(RowsAndColumns.class, shapeshifter);
+ Assertions.assertEquals(1, rac.numRows());
+ Assertions.assertEquals("b", rac.findColumn("dim").toAccessor().getObject(0));
+ }
+ }
+}
diff --git a/server/src/main/java/org/apache/druid/client/DirectDruidClient.java b/server/src/main/java/org/apache/druid/client/DirectDruidClient.java
index 2f963e0b270c..21addd635241 100644
--- a/server/src/main/java/org/apache/druid/client/DirectDruidClient.java
+++ b/server/src/main/java/org/apache/druid/client/DirectDruidClient.java
@@ -797,7 +797,8 @@ private void checkTotalBytesLimit(long bytes)
.setHeader(
HttpHeaders.Names.CONTENT_TYPE,
isSmile ? SmileMediaTypes.APPLICATION_JACKSON_SMILE : MediaType.APPLICATION_JSON
- ),
+ )
+ .setHeader(QueryResource.HEADER_NATIVE_QUERY_ROUTE, QueryResource.NATIVE_QUERY_ROUTE_LOCAL),
responseHandler,
Duration.millis(timeLeft)
);
@@ -880,8 +881,12 @@ private void cancelQuery(Query query, String cancelUrl)
try {
Future responseFuture = httpClient.go(
new Request(HttpMethod.DELETE, new URL(cancelUrl))
- .setContent(objectMapper.writeValueAsBytes(query))
- .setHeader(HttpHeaders.Names.CONTENT_TYPE, isSmile ? SmileMediaTypes.APPLICATION_JACKSON_SMILE : MediaType.APPLICATION_JSON),
+ .setContent(objectMapper.writeValueAsBytes(query))
+ .setHeader(
+ HttpHeaders.Names.CONTENT_TYPE,
+ isSmile ? SmileMediaTypes.APPLICATION_JACKSON_SMILE : MediaType.APPLICATION_JSON
+ )
+ .setHeader(QueryResource.HEADER_NATIVE_QUERY_ROUTE, QueryResource.NATIVE_QUERY_ROUTE_LOCAL),
StatusResponseHandler.getInstance(),
Duration.standardSeconds(1));
diff --git a/server/src/main/java/org/apache/druid/guice/BrokerQueryResourceModule.java b/server/src/main/java/org/apache/druid/guice/BrokerQueryResourceModule.java
new file mode 100644
index 000000000000..999e7547a631
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/guice/BrokerQueryResourceModule.java
@@ -0,0 +1,40 @@
+/*
+ * 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.guice;
+
+import com.google.inject.Binder;
+import org.apache.druid.initialization.DruidModule;
+import org.apache.druid.server.BrokerQueryResource;
+import org.apache.druid.server.ResponseContextConfig;
+import org.apache.druid.server.metrics.QueryCountStatsProvider;
+
+/** Registers the Broker's specialized native query HTTP resource. */
+public class BrokerQueryResourceModule implements DruidModule
+{
+ @Override
+ public void configure(final Binder binder)
+ {
+ binder.bind(ResponseContextConfig.class).toInstance(ResponseContextConfig.newConfig(false));
+ binder.bind(BrokerQueryResource.class).in(LazySingleton.class);
+ binder.bind(QueryCountStatsProvider.class).to(BrokerQueryResource.class).in(LazySingleton.class);
+ Jerseys.addResource(binder, BrokerQueryResource.class);
+ LifecycleModule.register(binder, BrokerQueryResource.class);
+ }
+}
diff --git a/server/src/main/java/org/apache/druid/guice/DruidBinders.java b/server/src/main/java/org/apache/druid/guice/DruidBinders.java
index 5c969d4ddf6b..7699a2b613ef 100644
--- a/server/src/main/java/org/apache/druid/guice/DruidBinders.java
+++ b/server/src/main/java/org/apache/druid/guice/DruidBinders.java
@@ -32,6 +32,7 @@
import org.apache.druid.query.QueryToolChest;
import org.apache.druid.segment.SegmentWrangler;
import org.apache.druid.segment.join.JoinableFactory;
+import org.apache.druid.server.DataSourceQueryHandler;
import java.lang.annotation.Annotation;
import java.util.Set;
@@ -147,6 +148,17 @@ public static MapBinder, SegmentWrangler> segmentWra
);
}
+ public static MapBinder, DataSourceQueryHandler> dataSourceQueryHandlerBinder(
+ Binder binder
+ )
+ {
+ return MapBinder.newMapBinder(
+ binder,
+ new TypeLiteral<>() {},
+ new TypeLiteral<>() {}
+ );
+ }
+
public static Multibinder joinableFactoryMultiBinder(Binder binder)
{
return Multibinder.newSetBinder(
diff --git a/server/src/main/java/org/apache/druid/guice/NativeQueryEngineModule.java b/server/src/main/java/org/apache/druid/guice/NativeQueryEngineModule.java
new file mode 100644
index 000000000000..cb161c344cd2
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/guice/NativeQueryEngineModule.java
@@ -0,0 +1,141 @@
+/*
+ * 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.guice;
+
+import com.google.common.collect.ImmutableSet;
+import com.google.inject.Binder;
+import com.google.inject.Module;
+import com.google.inject.multibindings.OptionalBinder;
+import com.google.inject.util.Modules;
+import org.apache.druid.initialization.DruidModule;
+import org.apache.druid.query.QuerySegmentWalker;
+import org.apache.druid.query.scan.ScanQuery;
+import org.apache.druid.server.NoopQuerySegmentWalker;
+import org.apache.druid.server.system.module.SystemTableModule;
+import org.apache.druid.server.system.module.SystemTableQueryResourceModule;
+
+import javax.annotation.Nullable;
+import java.util.List;
+
+/**
+ * Facade for the modules that provide native query execution infrastructure.
+ */
+public class NativeQueryEngineModule implements DruidModule
+{
+ private final Module executionModule;
+ private final Module queryResourceModule;
+
+ private NativeQueryEngineModule(
+ final Module executionModule,
+ final Module queryResourceModule
+ )
+ {
+ this.executionModule = executionModule;
+ this.queryResourceModule = queryResourceModule;
+ }
+
+ public static Builder builder()
+ {
+ return new Builder();
+ }
+
+ public static final class Builder
+ {
+ private Module overrideModule = binder -> {};
+ @Nullable
+ private Module queryResourceModule;
+ private boolean scanOnly;
+
+ private Builder()
+ {
+ }
+
+ /**
+ * Uses the minimum infrastructure required to serve Scan queries. This profile does not install processing or
+ * merge-buffer dependencies.
+ */
+ public Builder scanOnly()
+ {
+ scanOnly = true;
+ return this;
+ }
+
+ /** Applies role-specific overrides to the query execution bindings. */
+ public Builder withOverrideModule(final Module module)
+ {
+ overrideModule = module;
+ return this;
+ }
+
+ /** Replaces the standard native query HTTP resource with a role-specific resource module. */
+ public Builder withQueryResourceModule(final Module module)
+ {
+ queryResourceModule = module;
+ return this;
+ }
+
+ public NativeQueryEngineModule build()
+ {
+ final Module queryRunnerFactoryModule;
+ final Module querySegmentWalkerModule;
+ if (scanOnly) {
+ queryRunnerFactoryModule = new QueryRunnerFactoryModule(ImmutableSet.of(ScanQuery.class));
+ // Scan-only servers resolve node-local system tables without walking segments. Modules that need a real
+ // walker, such as SegmentSchemaCacheModule, replace this optional default with an explicit binding.
+ querySegmentWalkerModule = binder -> OptionalBinder.newOptionalBinder(binder, QuerySegmentWalker.class)
+ .setDefault()
+ .to(NoopQuerySegmentWalker.class)
+ .in(LazySingleton.class);
+ } else {
+ queryRunnerFactoryModule = new QueryRunnerFactoryModule();
+ querySegmentWalkerModule = binder -> {};
+ }
+
+ return new NativeQueryEngineModule(
+ Modules.override(
+ Modules.combine(
+ new QueryableModule(),
+ queryRunnerFactoryModule,
+ querySegmentWalkerModule
+ )
+ ).with(overrideModule),
+ queryResourceModule == null
+ ? scanOnly ? new SystemTableQueryResourceModule() : new QueryResourceModule()
+ : queryResourceModule
+ );
+ }
+ }
+
+ @Override
+ public void configure(final Binder binder)
+ {
+ binder.install(executionModule);
+ binder.install(new SegmentWranglerModule());
+ binder.install(new JoinableFactoryModule());
+ binder.install(new SystemTableModule());
+ binder.install(queryResourceModule);
+ }
+
+ @Override
+ public List getJacksonModules()
+ {
+ return new QueryableModule().getJacksonModules();
+ }
+}
diff --git a/services/src/main/java/org/apache/druid/guice/QueryablePeonModule.java b/server/src/main/java/org/apache/druid/guice/QueryResourceModule.java
similarity index 80%
rename from services/src/main/java/org/apache/druid/guice/QueryablePeonModule.java
rename to server/src/main/java/org/apache/druid/guice/QueryResourceModule.java
index 0d9c316b3ecf..3902946e0771 100644
--- a/services/src/main/java/org/apache/druid/guice/QueryablePeonModule.java
+++ b/server/src/main/java/org/apache/druid/guice/QueryResourceModule.java
@@ -22,13 +22,16 @@
import com.google.inject.Binder;
import org.apache.druid.initialization.DruidModule;
import org.apache.druid.server.QueryResource;
+import org.apache.druid.server.ResponseContextConfig;
import org.apache.druid.server.metrics.QueryCountStatsProvider;
-public class QueryablePeonModule implements DruidModule
+/** Registers the standard native query HTTP resource. */
+public class QueryResourceModule implements DruidModule
{
@Override
- public void configure(Binder binder)
+ public void configure(final Binder binder)
{
+ binder.bind(ResponseContextConfig.class).toInstance(ResponseContextConfig.newConfig(true));
binder.bind(QueryCountStatsProvider.class).to(QueryResource.class);
Jerseys.addResource(binder, QueryResource.class);
LifecycleModule.register(binder, QueryResource.class);
diff --git a/server/src/main/java/org/apache/druid/guice/QueryRunnerFactoryModule.java b/server/src/main/java/org/apache/druid/guice/QueryRunnerFactoryModule.java
index e991fac51b5d..0d7f2d2e6592 100644
--- a/server/src/main/java/org/apache/druid/guice/QueryRunnerFactoryModule.java
+++ b/server/src/main/java/org/apache/druid/guice/QueryRunnerFactoryModule.java
@@ -19,7 +19,9 @@
package org.apache.druid.guice;
+import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
import com.google.inject.Binder;
import com.google.inject.Key;
import com.google.inject.Provides;
@@ -52,6 +54,7 @@
import org.apache.druid.server.QuerySchedulerProvider;
import java.util.Map;
+import java.util.Set;
/**
*/
@@ -70,6 +73,32 @@ public class QueryRunnerFactoryModule extends QueryToolChestModule
.put(WindowOperatorQuery.class, WindowOperatorQueryQueryRunnerFactory.class)
.build();
+ private static final Set> SUPPORTED_QUERY_TYPES =
+ ImmutableSet.>builder()
+ .addAll(MAPPINGS.keySet())
+ .add(UnionQuery.class)
+ .build();
+
+ private final Set> queryTypes;
+
+ public QueryRunnerFactoryModule()
+ {
+ this(SUPPORTED_QUERY_TYPES);
+ }
+
+ public QueryRunnerFactoryModule(final Set extends Class extends Query>> queryTypes)
+ {
+ super(queryTypes);
+ Preconditions.checkArgument(
+ SUPPORTED_QUERY_TYPES.containsAll(queryTypes),
+ "Unsupported query types[%s]",
+ ImmutableSet.copyOf(queryTypes).stream()
+ .filter(queryType -> !SUPPORTED_QUERY_TYPES.contains(queryType))
+ .toList()
+ );
+ this.queryTypes = ImmutableSet.copyOf(queryTypes);
+ }
+
@Override
public void configure(Binder binder)
{
@@ -86,12 +115,18 @@ public void configure(Binder binder)
);
for (Map.Entry>, Class extends QueryRunnerFactory, ?>>> entry : MAPPINGS.entrySet()) {
- queryFactoryBinder.addBinding(entry.getKey()).to(entry.getValue());
- binder.bind(entry.getValue()).in(LazySingleton.class);
+ if (queryTypes.contains(entry.getKey())) {
+ queryFactoryBinder.addBinding(entry.getKey()).to(entry.getValue());
+ binder.bind(entry.getValue()).in(LazySingleton.class);
+ }
}
- DruidBinders.queryBinder(binder)
- .bindQueryLogic(UnionQuery.class, UnionQueryLogic.class);
+ if (queryTypes.contains(UnionQuery.class)) {
+ DruidBinders.queryBinder(binder)
+ .bindQueryLogic(UnionQuery.class, UnionQueryLogic.class);
+ } else {
+ DruidBinders.queryBinder(binder);
+ }
}
@LazySingleton
diff --git a/server/src/main/java/org/apache/druid/guice/QueryToolChestModule.java b/server/src/main/java/org/apache/druid/guice/QueryToolChestModule.java
index 72f0e142b297..d1ecdb6009fb 100644
--- a/server/src/main/java/org/apache/druid/guice/QueryToolChestModule.java
+++ b/server/src/main/java/org/apache/druid/guice/QueryToolChestModule.java
@@ -20,6 +20,7 @@
package org.apache.druid.guice;
import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
import com.google.inject.Binder;
import com.google.inject.Key;
import com.google.inject.Module;
@@ -64,6 +65,7 @@
import org.apache.druid.query.topn.TopNQueryQueryToolChest;
import java.util.Map;
+import java.util.Set;
/**
*/
@@ -75,7 +77,7 @@ public class QueryToolChestModule implements Module
public static final String TOPN_QUERY_METRICS_FACTORY_PROPERTY = "druid.query.topN.queryMetricsFactory";
public static final String SEARCH_QUERY_METRICS_FACTORY_PROPERTY = "druid.query.search.queryMetricsFactory";
- public final Map, Class extends QueryToolChest>> mappings =
+ private static final Map, Class extends QueryToolChest>> MAPPINGS =
ImmutableMap., Class extends QueryToolChest>>builder()
.put(DataSourceMetadataQuery.class, DataSourceQueryQueryToolChest.class)
.put(GroupByQuery.class, GroupByQueryQueryToolChest.class)
@@ -88,14 +90,30 @@ public class QueryToolChestModule implements Module
.put(WindowOperatorQuery.class, WindowOperatorQueryQueryToolChest.class)
.build();
+ public final Map, Class extends QueryToolChest>> mappings = MAPPINGS;
+
+ private final Set> queryTypes;
+
+ public QueryToolChestModule()
+ {
+ this(MAPPINGS.keySet());
+ }
+
+ public QueryToolChestModule(final Set extends Class extends Query>> queryTypes)
+ {
+ this.queryTypes = ImmutableSet.copyOf(queryTypes);
+ }
+
@Override
public void configure(Binder binder)
{
MapBinder, QueryToolChest> toolChests = DruidBinders.queryToolChestBinder(binder);
- for (Map.Entry, Class extends QueryToolChest>> entry : mappings.entrySet()) {
- toolChests.addBinding(entry.getKey()).to(entry.getValue());
- binder.bind(entry.getValue()).in(LazySingleton.class);
+ for (Map.Entry, Class extends QueryToolChest>> entry : MAPPINGS.entrySet()) {
+ if (queryTypes.contains(entry.getKey())) {
+ toolChests.addBinding(entry.getKey()).to(entry.getValue());
+ binder.bind(entry.getValue()).in(LazySingleton.class);
+ }
}
binder.bind(QueryToolChestWarehouse.class).to(ConglomerateBackedToolChestWarehouse.class);
@@ -103,11 +121,21 @@ public void configure(Binder binder)
JsonConfigProvider.bind(binder, "druid.query.default", DefaultQueryConfig.class);
// DefaultQueryContext defaults to the static DefaultQueryConfig; brokers override this binding.
binder.bind(QueryConfigProvider.class).to(DefaultQueryConfig.class);
- JsonConfigProvider.bind(binder, "druid.query.groupBy", GroupByQueryConfig.class);
- JsonConfigProvider.bind(binder, "druid.query.search", SearchQueryConfig.class);
- JsonConfigProvider.bind(binder, "druid.query.topN", TopNQueryConfig.class);
- JsonConfigProvider.bind(binder, "druid.query.segmentMetadata", SegmentMetadataQueryConfig.class);
- JsonConfigProvider.bind(binder, "druid.query.scan", ScanQueryConfig.class);
+ if (queryTypes.contains(GroupByQuery.class)) {
+ JsonConfigProvider.bind(binder, "druid.query.groupBy", GroupByQueryConfig.class);
+ }
+ if (queryTypes.contains(SearchQuery.class)) {
+ JsonConfigProvider.bind(binder, "druid.query.search", SearchQueryConfig.class);
+ }
+ if (queryTypes.contains(TopNQuery.class)) {
+ JsonConfigProvider.bind(binder, "druid.query.topN", TopNQueryConfig.class);
+ }
+ if (queryTypes.contains(SegmentMetadataQuery.class)) {
+ JsonConfigProvider.bind(binder, "druid.query.segmentMetadata", SegmentMetadataQueryConfig.class);
+ }
+ if (queryTypes.contains(ScanQuery.class)) {
+ JsonConfigProvider.bind(binder, "druid.query.scan", ScanQueryConfig.class);
+ }
PolyBind.createChoice(
binder,
@@ -120,48 +148,56 @@ public void configure(Binder binder)
.addBinding("default")
.to(DefaultGenericQueryMetricsFactory.class);
- PolyBind.createChoice(
- binder,
- GROUPBY_QUERY_METRICS_FACTORY_PROPERTY,
- Key.get(GroupByQueryMetricsFactory.class),
- Key.get(DefaultGroupByQueryMetricsFactory.class)
- );
- PolyBind
- .optionBinder(binder, Key.get(GroupByQueryMetricsFactory.class))
- .addBinding("default")
- .to(DefaultGroupByQueryMetricsFactory.class);
+ if (queryTypes.contains(GroupByQuery.class)) {
+ PolyBind.createChoice(
+ binder,
+ GROUPBY_QUERY_METRICS_FACTORY_PROPERTY,
+ Key.get(GroupByQueryMetricsFactory.class),
+ Key.get(DefaultGroupByQueryMetricsFactory.class)
+ );
+ PolyBind
+ .optionBinder(binder, Key.get(GroupByQueryMetricsFactory.class))
+ .addBinding("default")
+ .to(DefaultGroupByQueryMetricsFactory.class);
+ }
- PolyBind.createChoice(
- binder,
- TIMESERIES_QUERY_METRICS_FACTORY_PROPERTY,
- Key.get(TimeseriesQueryMetricsFactory.class),
- Key.get(DefaultTimeseriesQueryMetricsFactory.class)
- );
- PolyBind
- .optionBinder(binder, Key.get(TimeseriesQueryMetricsFactory.class))
- .addBinding("default")
- .to(DefaultTimeseriesQueryMetricsFactory.class);
+ if (queryTypes.contains(TimeseriesQuery.class)) {
+ PolyBind.createChoice(
+ binder,
+ TIMESERIES_QUERY_METRICS_FACTORY_PROPERTY,
+ Key.get(TimeseriesQueryMetricsFactory.class),
+ Key.get(DefaultTimeseriesQueryMetricsFactory.class)
+ );
+ PolyBind
+ .optionBinder(binder, Key.get(TimeseriesQueryMetricsFactory.class))
+ .addBinding("default")
+ .to(DefaultTimeseriesQueryMetricsFactory.class);
+ }
- PolyBind.createChoice(
- binder,
- TOPN_QUERY_METRICS_FACTORY_PROPERTY,
- Key.get(TopNQueryMetricsFactory.class),
- Key.get(DefaultTopNQueryMetricsFactory.class)
- );
- PolyBind
- .optionBinder(binder, Key.get(TopNQueryMetricsFactory.class))
- .addBinding("default")
- .to(DefaultTopNQueryMetricsFactory.class);
+ if (queryTypes.contains(TopNQuery.class)) {
+ PolyBind.createChoice(
+ binder,
+ TOPN_QUERY_METRICS_FACTORY_PROPERTY,
+ Key.get(TopNQueryMetricsFactory.class),
+ Key.get(DefaultTopNQueryMetricsFactory.class)
+ );
+ PolyBind
+ .optionBinder(binder, Key.get(TopNQueryMetricsFactory.class))
+ .addBinding("default")
+ .to(DefaultTopNQueryMetricsFactory.class);
+ }
- PolyBind.createChoice(
- binder,
- SEARCH_QUERY_METRICS_FACTORY_PROPERTY,
- Key.get(SearchQueryMetricsFactory.class),
- Key.get(DefaultSearchQueryMetricsFactory.class)
- );
- PolyBind
- .optionBinder(binder, Key.get(SearchQueryMetricsFactory.class))
- .addBinding("default")
- .to(DefaultSearchQueryMetricsFactory.class);
+ if (queryTypes.contains(SearchQuery.class)) {
+ PolyBind.createChoice(
+ binder,
+ SEARCH_QUERY_METRICS_FACTORY_PROPERTY,
+ Key.get(SearchQueryMetricsFactory.class),
+ Key.get(DefaultSearchQueryMetricsFactory.class)
+ );
+ PolyBind
+ .optionBinder(binder, Key.get(SearchQueryMetricsFactory.class))
+ .addBinding("default")
+ .to(DefaultSearchQueryMetricsFactory.class);
+ }
}
}
diff --git a/server/src/main/java/org/apache/druid/guice/QueryableModule.java b/server/src/main/java/org/apache/druid/guice/QueryableModule.java
index 7bd194132144..2a462734f7af 100644
--- a/server/src/main/java/org/apache/druid/guice/QueryableModule.java
+++ b/server/src/main/java/org/apache/druid/guice/QueryableModule.java
@@ -45,6 +45,7 @@ public class QueryableModule implements DruidModule
@Override
public void configure(Binder binder)
{
+ DruidBinders.dataSourceQueryHandlerBinder(binder);
binder.bind(RequestLogger.class).toProvider(RequestLoggerProvider.class).in(ManageLifecycle.class);
JsonConfigProvider.bindWithDefault(
binder,
diff --git a/server/src/main/java/org/apache/druid/server/DataSourceQueryHandler.java b/server/src/main/java/org/apache/druid/server/DataSourceQueryHandler.java
new file mode 100644
index 000000000000..6ad2cef59e3c
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/DataSourceQueryHandler.java
@@ -0,0 +1,44 @@
+/*
+ * 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.server;
+
+import org.apache.druid.query.Query;
+import org.apache.druid.query.QueryRunner;
+import org.apache.druid.server.security.AuthenticationResult;
+
+/**
+ * Creates a runner for a datasource whose execution is not provided by the normal segment walker. The query may
+ * contain the registered datasource below its root; the handler is responsible for resolving every matching vertex
+ * before delegating the remaining query to normal native execution.
+ */
+public interface DataSourceQueryHandler
+{
+ /**
+ * Creates the datasource runner for an authenticated request.
+ *
+ * @param executeLocally whether the request selected local execution with
+ * {@link QueryResource#HEADER_NATIVE_QUERY_ROUTE}
+ */
+ QueryRunner createRunner(
+ Query query,
+ AuthenticationResult authenticationResult,
+ boolean executeLocally
+ );
+}
diff --git a/server/src/main/java/org/apache/druid/server/QueryLifecycle.java b/server/src/main/java/org/apache/druid/server/QueryLifecycle.java
index 0f970a7fe0ce..ad7f8aac8500 100644
--- a/server/src/main/java/org/apache/druid/server/QueryLifecycle.java
+++ b/server/src/main/java/org/apache/druid/server/QueryLifecycle.java
@@ -34,6 +34,7 @@
import org.apache.druid.java.util.common.logger.Logger;
import org.apache.druid.java.util.emitter.service.ServiceEmitter;
import org.apache.druid.query.BaseQuery;
+import org.apache.druid.query.DataSource;
import org.apache.druid.query.DruidMetrics;
import org.apache.druid.query.GenericQueryMetricsFactory;
import org.apache.druid.query.Query;
@@ -41,6 +42,7 @@
import org.apache.druid.query.QueryInterruptedException;
import org.apache.druid.query.QueryMetrics;
import org.apache.druid.query.QueryPlus;
+import org.apache.druid.query.QueryRunner;
import org.apache.druid.query.QueryRunnerFactoryConglomerate;
import org.apache.druid.query.QuerySegmentWalker;
import org.apache.druid.query.QueryTimeoutException;
@@ -62,6 +64,7 @@
import javax.annotation.Nullable;
import javax.servlet.http.HttpServletRequest;
+import java.util.Collections;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
@@ -99,12 +102,14 @@ public class QueryLifecycle
private final AuthConfig authConfig;
private final PolicyEnforcer policyEnforcer;
private final QueryConfigSnapshot configSnapshot;
+ private final Map, DataSourceQueryHandler> dataSourceQueryHandlers;
private final long startMs;
private final long startNs;
private State state = State.NEW;
private AuthenticationResult authenticationResult;
private QueryToolChest toolChest;
+ private boolean executeNativeQueryLocally;
@MonotonicNonNull
private Query> baseQuery;
@@ -124,6 +129,7 @@ public QueryLifecycle(
final AuthConfig authConfig,
final PolicyEnforcer policyEnforcer,
final QueryConfigSnapshot configSnapshot,
+ final Map, DataSourceQueryHandler> dataSourceQueryHandlers,
final long startMs,
final long startNs
)
@@ -137,10 +143,41 @@ public QueryLifecycle(
this.authConfig = authConfig;
this.policyEnforcer = policyEnforcer;
this.configSnapshot = configSnapshot;
+ this.dataSourceQueryHandlers = dataSourceQueryHandlers;
this.startMs = startMs;
this.startNs = startNs;
}
+ public QueryLifecycle(
+ final QueryRunnerFactoryConglomerate conglomerate,
+ final QuerySegmentWalker texasRanger,
+ final GenericQueryMetricsFactory queryMetricsFactory,
+ final ServiceEmitter emitter,
+ final RequestLogger requestLogger,
+ final AuthorizerMapper authorizerMapper,
+ final AuthConfig authConfig,
+ final PolicyEnforcer policyEnforcer,
+ final QueryConfigSnapshot configSnapshot,
+ final long startMs,
+ final long startNs
+ )
+ {
+ this(
+ conglomerate,
+ texasRanger,
+ queryMetricsFactory,
+ emitter,
+ requestLogger,
+ authorizerMapper,
+ authConfig,
+ policyEnforcer,
+ configSnapshot,
+ Collections.emptyMap(),
+ startMs,
+ startNs
+ );
+ }
+
/**
* For callers who have already authorized their query, and where simplicity is desired over flexibility. This method
* does it all in one call. Logs and metrics are emitted when the Sequence is either fully iterated or throws an
@@ -271,6 +308,9 @@ public void initialize(final Query> baseQuery, @Nullable final Set cli
public AuthorizationResult authorize(HttpServletRequest req)
{
transition(State.INITIALIZED, State.AUTHORIZING);
+ executeNativeQueryLocally = QueryResource.NATIVE_QUERY_ROUTE_LOCAL.equals(
+ req.getHeader(QueryResource.HEADER_NATIVE_QUERY_ROUTE)
+ );
final Iterable resourcesToAuthorize = Iterables.concat(
Iterables.transform(
baseQuery.getDataSource().getTableNames(),
@@ -420,13 +460,49 @@ public QueryResponse execute()
final ResponseContext responseContext = DirectDruidClient.makeResponseContextForQuery();
@SuppressWarnings("unchecked")
- final Sequence res = QueryPlus.wrap((Query) baseQuery)
- .withIdentity(authenticationResult.getIdentity())
- .run(texasRanger, responseContext);
+ final Query query = (Query) baseQuery;
+ final DataSourceQueryHandler dataSourceQueryHandler = findDataSourceQueryHandler(query.getDataSource());
+ final QueryRunner queryRunner = dataSourceQueryHandler == null
+ ? query.getRunner(texasRanger)
+ : dataSourceQueryHandler.createRunner(
+ query,
+ authenticationResult,
+ executeNativeQueryLocally
+ );
+ final Sequence res = queryRunner.run(
+ QueryPlus.wrap(query).withIdentity(authenticationResult.getIdentity()),
+ responseContext
+ );
return new QueryResponse<>(res == null ? Sequences.empty() : res, responseContext);
}
+ /**
+ * Finds a handler for the root datasource or one of its descendants. A handler selected for a descendant receives
+ * the complete query so it can resolve the matching datasource vertices before normal native execution.
+ */
+ @Nullable
+ private DataSourceQueryHandler findDataSourceQueryHandler(final DataSource dataSource)
+ {
+ final DataSourceQueryHandler directHandler = dataSourceQueryHandlers.get(dataSource.getClass());
+ if (directHandler != null) {
+ return directHandler;
+ }
+
+ DataSourceQueryHandler descendantHandler = null;
+ for (final DataSource child : dataSource.getChildren()) {
+ final DataSourceQueryHandler childHandler = findDataSourceQueryHandler(child);
+ if (childHandler == null) {
+ continue;
+ }
+ if (descendantHandler != null && descendantHandler != childHandler) {
+ throw new ISE("Multiple datasource query handlers are required for datasource[%s]", dataSource);
+ }
+ descendantHandler = childHandler;
+ }
+ return descendantHandler;
+ }
+
/**
* Emits logs and metrics for this query.
*
diff --git a/server/src/main/java/org/apache/druid/server/QueryLifecycleFactory.java b/server/src/main/java/org/apache/druid/server/QueryLifecycleFactory.java
index d22d24ed3efb..2614301431b0 100644
--- a/server/src/main/java/org/apache/druid/server/QueryLifecycleFactory.java
+++ b/server/src/main/java/org/apache/druid/server/QueryLifecycleFactory.java
@@ -23,6 +23,7 @@
import org.apache.druid.client.BrokerViewOfBrokerConfig;
import org.apache.druid.guice.LazySingleton;
import org.apache.druid.java.util.emitter.service.ServiceEmitter;
+import org.apache.druid.query.DataSource;
import org.apache.druid.query.GenericQueryMetricsFactory;
import org.apache.druid.query.QueryConfigProvider;
import org.apache.druid.query.QueryRunnerFactoryConglomerate;
@@ -34,6 +35,7 @@
import org.apache.druid.server.security.AuthorizerMapper;
import javax.annotation.Nullable;
+import java.util.Map;
@LazySingleton
public class QueryLifecycleFactory
@@ -48,6 +50,7 @@ public class QueryLifecycleFactory
private final AuthConfig authConfig;
private final PolicyEnforcer policyEnforcer;
private final BrokerViewOfBrokerConfig brokerViewOfBrokerConfig;
+ private final Map, DataSourceQueryHandler> dataSourceQueryHandlers;
@Inject
public QueryLifecycleFactory(
@@ -60,6 +63,7 @@ public QueryLifecycleFactory(
final PolicyEnforcer policyEnforcer,
final AuthorizerMapper authorizerMapper,
final QueryConfigProvider queryConfigProvider,
+ final Map, DataSourceQueryHandler> dataSourceQueryHandlers,
@Nullable final BrokerViewOfBrokerConfig brokerViewOfBrokerConfig
)
{
@@ -70,11 +74,40 @@ public QueryLifecycleFactory(
this.requestLogger = requestLogger;
this.authorizerMapper = authorizerMapper;
this.queryConfigProvider = queryConfigProvider;
+ this.dataSourceQueryHandlers = dataSourceQueryHandlers;
this.authConfig = authConfig;
this.policyEnforcer = policyEnforcer;
this.brokerViewOfBrokerConfig = brokerViewOfBrokerConfig;
}
+ public QueryLifecycleFactory(
+ final QueryRunnerFactoryConglomerate conglomerate,
+ final QuerySegmentWalker texasRanger,
+ final GenericQueryMetricsFactory queryMetricsFactory,
+ final ServiceEmitter emitter,
+ final RequestLogger requestLogger,
+ final AuthConfig authConfig,
+ final PolicyEnforcer policyEnforcer,
+ final AuthorizerMapper authorizerMapper,
+ final QueryConfigProvider queryConfigProvider,
+ @Nullable final BrokerViewOfBrokerConfig brokerViewOfBrokerConfig
+ )
+ {
+ this(
+ conglomerate,
+ texasRanger,
+ queryMetricsFactory,
+ emitter,
+ requestLogger,
+ authConfig,
+ policyEnforcer,
+ authorizerMapper,
+ queryConfigProvider,
+ Map.of(),
+ brokerViewOfBrokerConfig
+ );
+ }
+
public QueryLifecycle factorize()
{
// Read once per query so the whole lifecycle sees one config, even if it is swapped mid-query.
@@ -93,6 +126,7 @@ public QueryLifecycle factorize()
authConfig,
policyEnforcer,
configSnapshot,
+ dataSourceQueryHandlers,
System.currentTimeMillis(),
System.nanoTime()
);
diff --git a/server/src/main/java/org/apache/druid/server/QueryResource.java b/server/src/main/java/org/apache/druid/server/QueryResource.java
index 97679d1edcf2..3213e81cf693 100644
--- a/server/src/main/java/org/apache/druid/server/QueryResource.java
+++ b/server/src/main/java/org/apache/druid/server/QueryResource.java
@@ -86,6 +86,9 @@ public class QueryResource implements QueryCountStatsProvider
public static final String HEADER_RESPONSE_CONTEXT = "X-Druid-Response-Context";
public static final String HEADER_IF_NONE_MATCH = "If-None-Match";
public static final String QUERY_ID_RESPONSE_HEADER = "X-Druid-Query-Id";
+ /** Selects where a native query sent through a forwarding service should execute. */
+ public static final String HEADER_NATIVE_QUERY_ROUTE = "X-Druid-Native-Query-Route";
+ public static final String NATIVE_QUERY_ROUTE_LOCAL = "local";
public static final String ERROR_MESSAGE_TRAILER_HEADER = "X-Error-Message";
public static final String RESPONSE_COMPLETE_TRAILER_HEADER = "X-Druid-Response-Complete";
public static final String HEADER_ETAG = "ETag";
@@ -287,7 +290,7 @@ public interface QueryMetricCounter
void incrementTimedOut();
}
- private Query> readQuery(
+ protected Query> readQuery(
final HttpServletRequest req,
final InputStream in,
final ResourceIOReaderWriterFactory.ResourceIOReaderWriter ioReaderWriter
diff --git a/server/src/main/java/org/apache/druid/server/http/CoordinatorRedirectInfo.java b/server/src/main/java/org/apache/druid/server/http/CoordinatorRedirectInfo.java
index 987ccb94842d..6e7d5c3c1a6e 100644
--- a/server/src/main/java/org/apache/druid/server/http/CoordinatorRedirectInfo.java
+++ b/server/src/main/java/org/apache/druid/server/http/CoordinatorRedirectInfo.java
@@ -22,6 +22,7 @@
import com.google.common.collect.ImmutableSet;
import com.google.inject.Inject;
import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.query.SystemTableDataSource;
import org.apache.druid.server.coordinator.DruidCoordinator;
import java.net.URL;
@@ -33,7 +34,9 @@ public class CoordinatorRedirectInfo implements RedirectInfo
{
private static final Set LOCAL_PATHS = ImmutableSet.of(
"/druid/coordinator/v1/leader",
- "/druid/coordinator/v1/isLeader"
+ "/druid/coordinator/v1/isLeader",
+ "/druid/v2",
+ "/druid/v2/"
);
private final DruidCoordinator coordinator;
@@ -47,7 +50,10 @@ public CoordinatorRedirectInfo(DruidCoordinator coordinator)
@Override
public boolean doLocal(String requestURI)
{
- return (requestURI != null && LOCAL_PATHS.contains(requestURI)) || coordinator.isLeader();
+ return (requestURI != null
+ && (LOCAL_PATHS.contains(requestURI)
+ || requestURI.startsWith("/druid/v2/" + SystemTableDataSource.NODE_QUERY_ID_PREFIX)))
+ || coordinator.isLeader();
}
@Override
diff --git a/server/src/main/java/org/apache/druid/server/system/SystemTableNotLeaderException.java b/server/src/main/java/org/apache/druid/server/system/SystemTableNotLeaderException.java
new file mode 100644
index 000000000000..b66b48f92767
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/system/SystemTableNotLeaderException.java
@@ -0,0 +1,38 @@
+/*
+ * 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.server.system;
+
+import org.apache.druid.java.util.common.ISE;
+
+/**
+ * Signals that a leader-only system table request reached a node that is no longer the leader.
+ *
+ *
The Broker may discover a leader and then contact it after leadership has changed. The native-query error
+ * response preserves this exception's class name, allowing {@code SystemTableQueryClient} to distinguish that race
+ * from an ordinary query failure, resolve the new leader, and retry once. A generic exception would not provide a
+ * reliable retry signal and retrying every failure could hide genuine query errors.
+ */
+public class SystemTableNotLeaderException extends ISE
+{
+ public SystemTableNotLeaderException(final String nodeRole)
+ {
+ super("Node role[%s] is not the current leader", nodeRole);
+ }
+}
diff --git a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableBrokerQueryHandler.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableBrokerQueryHandler.java
new file mode 100644
index 000000000000..1a8db1011306
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableBrokerQueryHandler.java
@@ -0,0 +1,63 @@
+/*
+ * 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.server.system.handler;
+
+import com.google.inject.Inject;
+import org.apache.druid.query.BadQueryContextException;
+import org.apache.druid.query.Query;
+import org.apache.druid.query.QueryRunner;
+import org.apache.druid.query.SystemTableDataSource;
+import org.apache.druid.server.DataSourceQueryHandler;
+import org.apache.druid.server.security.AuthenticationResult;
+
+/** Separates original Broker fanout from node-local requests received through the Broker's native endpoint. */
+public class SystemTableBrokerQueryHandler implements DataSourceQueryHandler
+{
+ private final SystemTableQueryClient systemTableQueryClient;
+ private final SystemTableQueryHandler localQueryHandler;
+
+ @Inject
+ public SystemTableBrokerQueryHandler(
+ final SystemTableQueryClient systemTableQueryClient,
+ final SystemTableQueryHandler localQueryHandler
+ )
+ {
+ this.systemTableQueryClient = systemTableQueryClient;
+ this.localQueryHandler = localQueryHandler;
+ }
+
+ @Override
+ public QueryRunner createRunner(
+ final Query query,
+ final AuthenticationResult authenticationResult,
+ final boolean executeLocally
+ )
+ {
+ if (executeLocally) {
+ if (!(query.getDataSource() instanceof SystemTableDataSource)) {
+ throw new BadQueryContextException(
+ "Local system-table execution requires a SystemTableDataSource at the query root"
+ );
+ }
+ return localQueryHandler.createRunner(query, authenticationResult, true);
+ }
+ return systemTableQueryClient.createRunner(query, authenticationResult, false);
+ }
+}
diff --git a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableLeaderLocator.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableLeaderLocator.java
new file mode 100644
index 000000000000..26f494c9c387
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableLeaderLocator.java
@@ -0,0 +1,31 @@
+/*
+ * 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.server.system.handler;
+
+import com.google.common.util.concurrent.ListenableFuture;
+
+import java.net.URI;
+
+/** Resolves the current leader for one Druid node role. */
+@FunctionalInterface
+public interface SystemTableLeaderLocator
+{
+ ListenableFuture findCurrentLeader();
+}
diff --git a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNode.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNode.java
new file mode 100644
index 000000000000..23f263eb8d83
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNode.java
@@ -0,0 +1,53 @@
+/*
+ * 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.server.system.handler;
+
+import org.apache.druid.discovery.DiscoveryDruidNode;
+import org.apache.druid.discovery.NodeRole;
+
+import java.util.LinkedHashSet;
+import java.util.Set;
+
+/** A physical Druid process selected to contribute rows to a native system-table query. */
+class SystemTableNode
+{
+ private final DiscoveryDruidNode discoveryNode;
+ private final Set nodeRoles = new LinkedHashSet<>();
+
+ SystemTableNode(final DiscoveryDruidNode discoveryNode)
+ {
+ this.discoveryNode = discoveryNode;
+ }
+
+ DiscoveryDruidNode getDiscoveryNode()
+ {
+ return discoveryNode;
+ }
+
+ Set getNodeRoles()
+ {
+ return nodeRoles;
+ }
+
+ void addNodeRole(final NodeRole nodeRole)
+ {
+ nodeRoles.add(nodeRole);
+ }
+}
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
new file mode 100644
index 000000000000..1706860c1407
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableNodeLocator.java
@@ -0,0 +1,150 @@
+/*
+ * 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.server.system.handler;
+
+import com.google.common.collect.Iterables;
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.inject.Inject;
+import org.apache.druid.client.DirectDruidClient;
+import org.apache.druid.client.coordinator.CoordinatorClient;
+import org.apache.druid.discovery.DiscoveryDruidNode;
+import org.apache.druid.discovery.DruidNodeDiscoveryProvider;
+import org.apache.druid.discovery.NodeRole;
+import org.apache.druid.java.util.common.ISE;
+import org.apache.druid.query.Query;
+import org.apache.druid.query.QueryInterruptedException;
+import org.apache.druid.query.QueryTimeoutException;
+import org.apache.druid.rpc.indexing.OverlordClient;
+import org.apache.druid.server.system.table.SystemTableDescriptor;
+import org.apache.druid.server.system.table.SystemTableRoutingMode;
+
+import java.net.URI;
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+/** Resolves a descriptor's logical routing policy to concrete Druid processes. */
+public class SystemTableNodeLocator
+{
+ private final DruidNodeDiscoveryProvider discoveryProvider;
+ private final Map leaderLocators;
+
+ @Inject
+ public SystemTableNodeLocator(
+ final DruidNodeDiscoveryProvider discoveryProvider,
+ final CoordinatorClient coordinatorClient,
+ final OverlordClient overlordClient
+ )
+ {
+ this.discoveryProvider = discoveryProvider;
+ this.leaderLocators = Map.of(
+ NodeRole.COORDINATOR,
+ coordinatorClient::findCurrentLeader,
+ NodeRole.OVERLORD,
+ overlordClient::findCurrentLeader
+ );
+ }
+
+ List locate(final SystemTableDescriptor descriptor, final Query> query)
+ {
+ if (descriptor.getRoutingMode() == SystemTableRoutingMode.ALL_NODES) {
+ return discoverAllNodes(descriptor);
+ }
+
+ final NodeRole leaderRole = Iterables.getOnlyElement(descriptor.getNodeRoles());
+ final URI leaderUri = findLeader(leaderRole, query);
+ // Resolve leadership before taking the discovery snapshot. A leader election may complete while discovery still
+ // contains the previous membership, so taking the snapshot first can reject a valid newly elected leader.
+ return discoverAllNodes(descriptor).stream()
+ .filter(node -> sameServer(leaderUri, node.getDiscoveryNode().getDruidNode().getUriToUse()))
+ .findFirst()
+ .map(List::of)
+ .orElseThrow(
+ () -> new ISE(
+ "Current leader[%s] for role[%s] is not present in service discovery",
+ leaderUri,
+ leaderRole
+ )
+ );
+ }
+
+ private List discoverAllNodes(final SystemTableDescriptor descriptor)
+ {
+ final Map nodes = new LinkedHashMap<>();
+ for (final NodeRole nodeRole : descriptor.getNodeRoles()) {
+ for (final DiscoveryDruidNode node : discoveryProvider.getForNodeRole(nodeRole).getAllNodes()) {
+ nodes.computeIfAbsent(
+ node.getDruidNode().getHostAndPortToUse(),
+ ignored -> new SystemTableNode(node)
+ ).addNodeRole(node.getNodeRole());
+ }
+ }
+ return new ArrayList<>(nodes.values());
+ }
+
+ private URI findLeader(final NodeRole nodeRole, final Query> query)
+ {
+ final SystemTableLeaderLocator leaderLocator = leaderLocators.get(nodeRole);
+ if (leaderLocator == null) {
+ throw new ISE("Leader-only system-table routing is not supported for role[%s]", nodeRole);
+ }
+
+ final ListenableFuture leaderFuture = leaderLocator.findCurrentLeader();
+ final long timeLeft = query.context().getLong(DirectDruidClient.QUERY_FAIL_TIME) - System.currentTimeMillis();
+ if (timeLeft <= 0) {
+ leaderFuture.cancel(true);
+ throw new QueryTimeoutException("Timed out while locating the current leader for role[" + nodeRole + "]");
+ }
+ try {
+ return leaderFuture.get(timeLeft, TimeUnit.MILLISECONDS);
+ }
+ catch (InterruptedException e) {
+ leaderFuture.cancel(true);
+ Thread.currentThread().interrupt();
+ throw QueryInterruptedException.wrapIfNeeded(e);
+ }
+ catch (TimeoutException e) {
+ leaderFuture.cancel(true);
+ throw new QueryTimeoutException("Timed out while locating the current leader for role[" + nodeRole + "]");
+ }
+ catch (ExecutionException e) {
+ throw QueryInterruptedException.wrapIfNeeded(e.getCause() == null ? e : e.getCause());
+ }
+ }
+
+ static boolean sameServer(final URI first, final URI second)
+ {
+ return first.getScheme().equalsIgnoreCase(second.getScheme())
+ && first.getHost().equalsIgnoreCase(second.getHost())
+ && effectivePort(first) == effectivePort(second);
+ }
+
+ private static int effectivePort(final URI uri)
+ {
+ if (uri.getPort() >= 0) {
+ return uri.getPort();
+ }
+ return "https".equalsIgnoreCase(uri.getScheme()) ? 443 : 80;
+ }
+}
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
new file mode 100644
index 000000000000..11d4121448dd
--- /dev/null
+++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java
@@ -0,0 +1,901 @@
+/*
+ * 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.server.system.handler;
+
+import com.google.common.base.Preconditions;
+import com.google.common.collect.Iterables;
+import com.google.inject.Inject;
+import io.netty.channel.ChannelException;
+import org.apache.druid.client.DirectDruidClient;
+import org.apache.druid.client.DirectDruidClientFactory;
+import org.apache.druid.client.DruidServer;
+import org.apache.druid.discovery.DiscoveryDruidNode;
+import org.apache.druid.guice.annotations.Self;
+import org.apache.druid.java.util.common.ISE;
+import org.apache.druid.java.util.common.JodaUtils;
+import org.apache.druid.java.util.common.guava.BaseSequence;
+import org.apache.druid.java.util.common.guava.LazySequence;
+import org.apache.druid.java.util.common.guava.Sequence;
+import org.apache.druid.java.util.common.guava.Sequences;
+import org.apache.druid.java.util.common.guava.Yielder;
+import org.apache.druid.java.util.common.guava.Yielders;
+import org.apache.druid.query.BaseQuery;
+import org.apache.druid.query.DataSource;
+import org.apache.druid.query.Druids;
+import org.apache.druid.query.InlineDataSource;
+import org.apache.druid.query.Query;
+import org.apache.druid.query.QueryContexts;
+import org.apache.druid.query.QueryDataSource;
+import org.apache.druid.query.QueryException;
+import org.apache.druid.query.QueryInterruptedException;
+import org.apache.druid.query.QueryPlus;
+import org.apache.druid.query.QueryRunner;
+import org.apache.druid.query.QuerySegmentWalker;
+import org.apache.druid.query.QueryTimeoutException;
+import org.apache.druid.query.SystemTableDataSource;
+import org.apache.druid.query.context.ResponseContext;
+import org.apache.druid.query.filter.AndDimFilter;
+import org.apache.druid.query.filter.DimFilter;
+import org.apache.druid.query.operator.OperatorFactory;
+import org.apache.druid.query.operator.ScanOperatorFactory;
+import org.apache.druid.query.operator.WindowOperatorQuery;
+import org.apache.druid.query.scan.ScanQuery;
+import org.apache.druid.query.scan.ScanResultValue;
+import org.apache.druid.segment.VirtualColumn;
+import org.apache.druid.segment.VirtualColumns;
+import org.apache.druid.server.DataSourceQueryHandler;
+import org.apache.druid.server.DruidNode;
+import org.apache.druid.server.QueryScheduler;
+import org.apache.druid.server.coordination.ServerType;
+import org.apache.druid.server.security.AuthenticationResult;
+import org.apache.druid.server.security.AuthorizerMapper;
+import org.apache.druid.server.security.Escalator;
+import org.apache.druid.server.system.SystemTableNotLeaderException;
+import org.apache.druid.server.system.table.SystemTableDescriptor;
+import org.apache.druid.server.system.table.SystemTableRoutingMode;
+import org.apache.druid.utils.CloseableUtils;
+
+import javax.annotation.Nullable;
+import java.io.Closeable;
+import java.io.EOFException;
+import java.io.IOException;
+import java.net.SocketException;
+import java.net.UnknownHostException;
+import java.nio.channels.ClosedChannelException;
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.NoSuchElementException;
+import java.util.Optional;
+import java.util.UUID;
+import java.util.function.Supplier;
+
+/**
+ * Recursively replaces native system-table datasources with authorized inline rows fetched from nodes, then
+ * executes the rebuilt query through the Broker's normal native-query walker.
+ */
+public class SystemTableQueryClient implements DataSourceQueryHandler
+{
+ private static final String SYSTEM_TABLE_TIER = "_system";
+
+ private final SystemTableNodeLocator nodeLocator;
+ private final DirectDruidClientFactory directDruidClientFactory;
+ private final QueryScheduler queryScheduler;
+ private final QuerySegmentWalker querySegmentWalker;
+ private final Map tableDescriptors;
+ private final AuthorizerMapper authorizerMapper;
+ private final SystemTableQueryHandler localQueryHandler;
+ private final AuthenticationResult escalatedAuthenticationResult;
+ private final DruidNode selfNode;
+
+ @Inject
+ public SystemTableQueryClient(
+ final SystemTableNodeLocator nodeLocator,
+ final DirectDruidClientFactory directDruidClientFactory,
+ final QueryScheduler queryScheduler,
+ final QuerySegmentWalker querySegmentWalker,
+ final Map tableDescriptors,
+ final AuthorizerMapper authorizerMapper,
+ final SystemTableQueryHandler localQueryHandler,
+ final Escalator escalator,
+ @Self final DruidNode selfNode
+ )
+ {
+ this.nodeLocator = nodeLocator;
+ this.directDruidClientFactory = directDruidClientFactory;
+ this.queryScheduler = queryScheduler;
+ this.querySegmentWalker = querySegmentWalker;
+ this.tableDescriptors = tableDescriptors;
+ this.authorizerMapper = authorizerMapper;
+ this.localQueryHandler = localQueryHandler;
+ this.escalatedAuthenticationResult = escalator.createEscalatedAuthenticationResult();
+ this.selfNode = selfNode;
+ }
+
+ @Override
+ public QueryRunner createRunner(
+ final Query query,
+ final AuthenticationResult authenticationResult,
+ final boolean executeLocally
+ )
+ {
+ return (queryPlus, responseContext) -> {
+ final NodeSequencesCloser nodeSequencesCloser = new NodeSequencesCloser(queryScheduler);
+ try {
+ final DataSource resolvedDataSource = resolveDataSource(
+ query.getDataSource(),
+ query,
+ authenticationResult,
+ responseContext,
+ nodeSequencesCloser
+ );
+ final Query resolvedQuery = query.withDataSource(resolvedDataSource).withOverriddenContext(
+ Map.of(QueryContexts.QUERY_RESOURCE_ID, UUID.randomUUID().toString())
+ );
+ final QueryRunner queryRunner = querySegmentWalker.getQueryRunnerForIntervals(
+ resolvedQuery,
+ resolvedQuery.getIntervals()
+ );
+ return Sequences.withBaggage(
+ queryRunner.run(queryPlus.withQuery(resolvedQuery), responseContext),
+ nodeSequencesCloser
+ );
+ }
+ catch (Throwable t) {
+ throw CloseableUtils.closeAndWrapInCatch(t, nodeSequencesCloser);
+ }
+ };
+ }
+
+ private DataSource resolveDataSource(
+ final DataSource dataSource,
+ final Query> owningQuery,
+ final AuthenticationResult authenticationResult,
+ final ResponseContext responseContext,
+ final NodeSequencesCloser nodeSequencesCloser
+ )
+ {
+ if (dataSource instanceof SystemTableDataSource systemTableDataSource) {
+ return resolveSystemTableDataSource(
+ systemTableDataSource,
+ owningQuery,
+ authenticationResult,
+ responseContext,
+ nodeSequencesCloser
+ );
+ }
+
+ final Query> childOwningQuery = dataSource instanceof QueryDataSource queryDataSource
+ ? queryDataSource.getQuery()
+ : owningQuery;
+ final List resolvedChildren = new ArrayList<>();
+ for (final DataSource child : dataSource.getChildren()) {
+ resolvedChildren.add(
+ resolveDataSource(
+ child,
+ childOwningQuery,
+ authenticationResult,
+ responseContext,
+ nodeSequencesCloser
+ )
+ );
+ }
+ return dataSource.withChildren(resolvedChildren);
+ }
+
+ private InlineDataSource resolveSystemTableDataSource(
+ final SystemTableDataSource dataSource,
+ final Query> owningQuery,
+ final AuthenticationResult authenticationResult,
+ final ResponseContext responseContext,
+ final NodeSequencesCloser nodeSequencesCloser
+ )
+ {
+ final SystemTableDescriptor descriptor = tableDescriptors.get(dataSource.getTable());
+ if (descriptor == null) {
+ throw new ISE("No routing descriptor is registered for system table[%s]", dataSource.getTable());
+ }
+
+ final ScanQuery nodeQuery = makeNodeQuery(dataSource, descriptor, owningQuery);
+ final List> nodeRunners = makeNodeRunners(
+ nodeQuery,
+ descriptor,
+ nodeSequencesCloser
+ );
+ if (nodeRunners.isEmpty() && !descriptor.isEmptyDiscoveryAllowed()) {
+ throw new ISE("No node is available to serve system table[%s]", dataSource.getTable());
+ }
+
+ /*
+ * Keep the scan result transport lazy. The node response is a Sequence, so this lets the local query
+ * runner consume and process rows while the HTTP response is still arriving instead of materializing every
+ * node row in two Broker-side lists. The scan transport is intentionally kept behind this adapter so a
+ * future node-side aggregation mode can use a different result adapter without changing routing or
+ * authorization orchestration here.
+ */
+ final NodeSequenceCloser nodeSequenceCloser = new NodeSequenceCloser();
+ nodeSequencesCloser.add(nodeSequenceCloser);
+ final Iterable