This is an automated email from the ASF dual-hosted git repository. FrankChen021 pushed a commit to branch codex/native-sys-segments in repository https://gitbox.apache.org/repos/asf/druid.git
commit 4d113ce1fb0872bf9fc478e00bffb575afeeebd9 Author: Frank Chen <[email protected]> AuthorDate: Mon Aug 31 16:40:03 2026 +0800 perf(sql): optimize local system table queries --- .../system/handler/SystemTableQueryClient.java | 55 +++++------ .../system/handler/SystemTableQueryHandler.java | 102 +++++++++++++++++---- .../system/table/SystemTableDataProvider.java | 27 ++++++ .../server/system/SystemTableQueryHandlerTest.java | 96 +++++++++++++++++++ .../system/handler/SystemTableQueryClientTest.java | 46 +++++----- .../calcite/schema/SegmentsTableDataProvider.java | 23 ++++- .../schema/SegmentsTableDataProviderTest.java | 22 +++++ 7 files changed, 299 insertions(+), 72 deletions(-) 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 58577802cbb..44e2a0db525 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 @@ -216,6 +216,9 @@ public class SystemTableQueryClient implements DataSourceQueryHandler } final ScanQuery nodeQuery = makeNodeQuery(dataSource, descriptor, owningQuery); + if (descriptor.getRoutingMode() == SystemTableRoutingMode.LOCAL_ONLY) { + return localQueryHandler.resolveDataSource(nodeQuery, authenticationResult); + } final List<QueryRunner<ScanResultValue>> nodeRunners = makeNodeRunners( nodeQuery, descriptor, @@ -270,7 +273,7 @@ public class SystemTableQueryClient implements DataSourceQueryHandler .limit(Long.MAX_VALUE) .filters(nodeFilter) .virtualColumns(nodeVirtualColumns(dataSource, owningQuery, nodeFilter)) - .columns(descriptor.getRowSignature()) + .columns(nodeColumns(dataSource, descriptor, owningQuery)) .context(nodeContext) .build(); } @@ -282,11 +285,6 @@ public class SystemTableQueryClient implements DataSourceQueryHandler ) { final List<QueryRunner<ScanResultValue>> nodeRunners = new ArrayList<>(); - if (descriptor.getRoutingMode() == SystemTableRoutingMode.LOCAL_ONLY) { - nodeRunners.add((queryPlus, responseContext) -> runLocalNodeQuery(nodeQuery, queryPlus, responseContext)); - return nodeRunners; - } - for (final SystemTableNode node : nodeLocator.locate(descriptor, nodeQuery)) { nodeRunners.add( (queryPlus, responseContext) -> recoverNodeFailure( @@ -310,30 +308,6 @@ 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 - ) - ); - final QueryRunner<ScanResultValue> localRunner = localQueryHandler.createRunner( - subNativeQuery, - escalatedAuthenticationResult, - true - ); - return localRunner.run(queryPlus.withQuery(subNativeQuery), responseContext); - } - private Sequence<ScanResultValue> runNodeQuery( final ScanQuery nodeQuery, final SystemTableNode node, @@ -861,6 +835,27 @@ public class SystemTableQueryClient implements DataSourceQueryHandler { } + private static List<String> nodeColumns( + final SystemTableDataSource dataSource, + final SystemTableDescriptor descriptor, + final Query<?> query + ) + { + if (descriptor.getRoutingMode() != SystemTableRoutingMode.LOCAL_ONLY + || !dataSource.equals(query.getDataSource()) + || query.getRequiredColumns() == null) { + return descriptor.getRowSignature().getColumnNames(); + } + + final List<String> columns = descriptor.getRowSignature() + .getColumnNames() + .stream() + .filter(query.getRequiredColumns()::contains) + .toList(); + // A count-only query still needs one physical column so the inline cursor retains the source row cardinality. + return columns.isEmpty() ? List.of(descriptor.getRowSignature().getColumnName(0)) : columns; + } + @Nullable private static DimFilter nodeFilter(final SystemTableDataSource dataSource, final Query<?> query) { diff --git a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryHandler.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryHandler.java index ae117467a9a..311621999c3 100644 --- a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryHandler.java +++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryHandler.java @@ -19,6 +19,7 @@ package org.apache.druid.server.system.handler; +import com.google.common.collect.Iterables; import com.google.inject.Inject; import org.apache.druid.client.DirectDruidClient; import org.apache.druid.java.util.common.IAE; @@ -36,6 +37,7 @@ import org.apache.druid.query.scan.ScanQuery; import org.apache.druid.query.scan.ScanQueryEngine; import org.apache.druid.segment.InlineSegmentWrangler; import org.apache.druid.segment.Segment; +import org.apache.druid.segment.column.RowSignature; import org.apache.druid.server.DataSourceQueryHandler; import org.apache.druid.server.security.AuthenticationResult; import org.apache.druid.server.security.AuthorizerMapper; @@ -43,9 +45,10 @@ import org.apache.druid.server.system.table.SystemTableDataProvider; import org.apache.druid.server.system.table.SystemTableDescriptor; import org.apache.druid.server.system.table.SystemTablePushdownFilter; +import java.util.List; import java.util.Map; -/** Resolves one node-local system table and returns its rows through the standard native Scan stack. */ +/** Resolves one node-local system table either to inline rows or through the standard native Scan stack. */ public class SystemTableQueryHandler implements DataSourceQueryHandler { private final Map<String, SystemTableDataProvider> dataSuppliers; @@ -79,29 +82,24 @@ public class SystemTableQueryHandler implements DataSourceQueryHandler } final SystemTableDataSource dataSource = (SystemTableDataSource) query.getDataSource(); - final SystemTableDataProvider dataSupplier = dataSuppliers.get(dataSource.getTable()); - if (dataSupplier == null) { - throw new ISE("System table[%s] is not served by this node", dataSource.getTable()); - } - final SystemTableDescriptor descriptor = tableDescriptors.get(dataSource.getTable()); - if (descriptor == null) { - throw new ISE("No descriptor is registered for system table[%s]", dataSource.getTable()); - } + final SystemTableDataProvider dataSupplier = getDataSupplier(dataSource); + final SystemTableDescriptor descriptor = getDescriptor(dataSource); return (queryPlus, responseContext) -> { - final Iterable<Object[]> suppliedRows = () -> dataSupplier.getRows( - SystemTablePushdownFilter.extract(query, dataSupplier.getPushdownFilters()), + final Iterable<Object[]> authorizedRows = getAuthorizedRows( + dataSupplier, + descriptor, + (ScanQuery) query, requestAuthenticationResult - ).iterator(); - final Iterable<Object[]> authorizedRows = descriptor.getRowAuthorizer().filterAuthorizedRows( - suppliedRows, - requestAuthenticationResult, - authorizerMapper + ); + final Iterable<Object[]> projectedRows = Iterables.transform( + authorizedRows, + row -> dataSupplier.projectRow(row, null) ); final ScanQuery resolvedQuery = Druids.ScanQueryBuilder.copy((ScanQuery) query) .dataSource( InlineDataSource.fromIterable( - authorizedRows, + projectedRows, descriptor.getRowSignature() ) ) @@ -121,6 +119,76 @@ public class SystemTableQueryHandler implements DataSourceQueryHandler }; } + /** Resolves a local system table directly to lazily projected, user-authorized inline rows. */ + public InlineDataSource resolveDataSource( + final ScanQuery query, + final AuthenticationResult requestAuthenticationResult + ) + { + final SystemTableDataSource dataSource = (SystemTableDataSource) query.getDataSource(); + final SystemTableDataProvider dataSupplier = getDataSupplier(dataSource); + final SystemTableDescriptor descriptor = getDescriptor(dataSource); + final List<String> columns = query.getColumns(); + final int[] projects = new int[columns.size()]; + final RowSignature.Builder signatureBuilder = RowSignature.builder(); + for (int i = 0; i < columns.size(); i++) { + final String column = columns.get(i); + final int columnNumber = descriptor.getRowSignature().indexOf(column); + if (columnNumber < 0) { + throw new IAE("Column[%s] is not present in system table[%s]", column, dataSource.getTable()); + } + projects[i] = columnNumber; + signatureBuilder.add(column, descriptor.getRowSignature().getColumnType(columnNumber).orElse(null)); + } + + final Iterable<Object[]> authorizedRows = getAuthorizedRows( + dataSupplier, + descriptor, + query, + requestAuthenticationResult + ); + return InlineDataSource.fromIterable( + Iterables.transform(authorizedRows, row -> dataSupplier.projectRow(row, projects)), + signatureBuilder.build() + ); + } + + private Iterable<Object[]> getAuthorizedRows( + final SystemTableDataProvider dataSupplier, + final SystemTableDescriptor descriptor, + final ScanQuery query, + final AuthenticationResult requestAuthenticationResult + ) + { + final Iterable<Object[]> suppliedRows = () -> dataSupplier.getRawRows( + SystemTablePushdownFilter.extract(query, dataSupplier.getPushdownFilters()), + requestAuthenticationResult + ).iterator(); + return descriptor.getRowAuthorizer().filterAuthorizedRows( + suppliedRows, + requestAuthenticationResult, + authorizerMapper + ); + } + + private SystemTableDataProvider getDataSupplier(final SystemTableDataSource dataSource) + { + final SystemTableDataProvider dataSupplier = dataSuppliers.get(dataSource.getTable()); + if (dataSupplier == null) { + throw new ISE("System table[%s] is not served by this node", dataSource.getTable()); + } + return dataSupplier; + } + + private SystemTableDescriptor getDescriptor(final SystemTableDataSource dataSource) + { + final SystemTableDescriptor descriptor = tableDescriptors.get(dataSource.getTable()); + if (descriptor == null) { + throw new ISE("No descriptor is registered for system table[%s]", dataSource.getTable()); + } + return descriptor; + } + @SuppressWarnings("unchecked") private static <T> Sequence<T> runScan( final ScanQueryEngine scanQueryEngine, diff --git a/server/src/main/java/org/apache/druid/server/system/table/SystemTableDataProvider.java b/server/src/main/java/org/apache/druid/server/system/table/SystemTableDataProvider.java index 7f40f001716..a4d6aa18943 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/SystemTableDataProvider.java +++ b/server/src/main/java/org/apache/druid/server/system/table/SystemTableDataProvider.java @@ -23,6 +23,7 @@ import jakarta.validation.constraints.NotNull; import org.apache.druid.query.filter.DimFilter; import org.apache.druid.server.security.AuthenticationResult; +import javax.annotation.Nullable; import java.util.Collections; import java.util.List; @@ -42,4 +43,30 @@ public interface SystemTableDataProvider @NotNull List<DimFilter> filters, AuthenticationResult internalAuthenticationResult ); + + /** + * Returns full-width rows before any type conversion that can be deferred until column projection. Providers whose + * rows already match the descriptor signature can use the default implementation. + */ + default Iterable<Object[]> getRawRows( + @NotNull final List<DimFilter> filters, + final AuthenticationResult internalAuthenticationResult + ) + { + return getRows(filters, internalAuthenticationResult); + } + + /** Projects one authorized raw row into the columns requested by the native query. */ + default Object[] projectRow(final Object[] row, @Nullable final int[] projects) + { + if (projects == null) { + return row; + } + + final Object[] projectedRow = new Object[projects.length]; + for (int i = 0; i < projects.length; i++) { + projectedRow[i] = row[projects[i]]; + } + return projectedRow; + } } diff --git a/server/src/test/java/org/apache/druid/server/system/SystemTableQueryHandlerTest.java b/server/src/test/java/org/apache/druid/server/system/SystemTableQueryHandlerTest.java index 33834a1e10f..866f2c86964 100644 --- a/server/src/test/java/org/apache/druid/server/system/SystemTableQueryHandlerTest.java +++ b/server/src/test/java/org/apache/druid/server/system/SystemTableQueryHandlerTest.java @@ -21,6 +21,7 @@ package org.apache.druid.server.system; import org.apache.druid.discovery.NodeRole; import org.apache.druid.query.Druids; +import org.apache.druid.query.InlineDataSource; import org.apache.druid.query.QueryPlus; import org.apache.druid.query.QueryRunner; import org.apache.druid.query.SystemTableDataSource; @@ -146,4 +147,99 @@ public class SystemTableQueryHandlerTest ); } + @Test + public void testDirectResolutionAuthorizesRawRowsBeforeProjection() + { + final AtomicInteger getRawRowsCalls = new AtomicInteger(); + final AtomicInteger projectRowCalls = new AtomicInteger(); + final AtomicInteger authorizationCalls = new AtomicInteger(); + final SystemTableDataProvider supplier = new SystemTableDataProvider() + { + @Override + public Iterable<Object[]> getRows( + final List<DimFilter> filters, + final AuthenticationResult authenticationResult + ) + { + throw new AssertionError("Direct resolution must use raw rows"); + } + + @Override + public Iterable<Object[]> getRawRows( + final List<DimFilter> filters, + final AuthenticationResult authenticationResult + ) + { + getRawRowsCalls.incrementAndGet(); + return List.of( + new Object[]{"task-a", 10L}, + new Object[]{"task-b", 20L} + ); + } + + @Override + public Object[] projectRow(final Object[] row, final int[] projects) + { + projectRowCalls.incrementAndGet(); + return SystemTableDataProvider.super.projectRow(row, projects); + } + }; + final SystemTableDescriptor descriptor = new SystemTableDescriptor() + { + @Override + public String getTableName() + { + return "test"; + } + + @Override + public Set<NodeRole> getNodeRoles() + { + return Set.of(); + } + + @Override + public RowSignature getRowSignature() + { + return ROW_SIGNATURE; + } + + @Override + public SystemTableRowAuthorizer getRowAuthorizer() + { + return (rows, authenticationResult, authorizerMapper) -> { + authorizationCalls.incrementAndGet(); + return () -> java.util.stream.StreamSupport.stream(rows.spliterator(), false) + .filter(row -> "task-b".equals(row[0])) + .iterator(); + }; + } + }; + final SystemTableQueryHandler handler = new SystemTableQueryHandler( + Map.of("test", supplier), + Map.of("test", descriptor), + new ScanQueryEngine(), + new AuthorizerMapper(Map.of()) + ); + final ScanQuery query = Druids.newScanQueryBuilder() + .dataSource(new SystemTableDataSource("test")) + .eternityInterval() + .columns(List.of("duration")) + .build(); + + final InlineDataSource inlineDataSource = handler.resolveDataSource( + query, + new AuthenticationResult("alice", "allow", "external", null) + ); + + Assertions.assertEquals(List.of("duration"), inlineDataSource.getRowSignature().getColumnNames()); + Assertions.assertEquals(0, getRawRowsCalls.get()); + Assertions.assertEquals(0, projectRowCalls.get()); + Assertions.assertEquals(1, authorizationCalls.get()); + final List<Object[]> rows = inlineDataSource.getRowsAsList(); + Assertions.assertEquals(1, getRawRowsCalls.get()); + Assertions.assertEquals(1, projectRowCalls.get()); + Assertions.assertArrayEquals(new Object[]{20L}, rows.get(0)); + } + } 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 45e3bfafae8..402e21b2d3e 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 @@ -89,6 +89,7 @@ import org.joda.time.Duration; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; +import org.mockito.ArgumentCaptor; import org.mockito.ArgumentMatchers; import org.mockito.Mockito; @@ -349,24 +350,19 @@ public class SystemTableQueryClientTest 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( - new ScanResultValue( - null, - descriptor.getRowSignature().getColumnNames(), - List.of((Object) new Object[descriptor.getRowSignature().size()]) - ) - ) - ) - ).when(localQueryHandler).createRunner( - ArgumentMatchers.any(), - Mockito.same(escalatedAuthenticationResult), - Mockito.eq(true) + Mockito.when(escalator.createEscalatedAuthenticationResult()).thenReturn( + new AuthenticationResult("system", "allow", "system", null) + ); + Mockito.when(localQueryHandler.resolveDataSource(ArgumentMatchers.any(), Mockito.same(AUTHENTICATION_RESULT))) + .thenReturn( + InlineDataSource.fromIterable( + Collections.singletonList(new Object[]{"foo", 1L}), + RowSignature.builder() + .add("datasource", ColumnType.STRING) + .add("size", ColumnType.LONG) + .build() + ) ); final QuerySegmentWalker querySegmentWalker = Mockito.mock(QuerySegmentWalker.class); @@ -383,7 +379,12 @@ public class SystemTableQueryClientTest escalator, nonMatchingSelfNode() ); - final ScanQuery query = query(descriptor, Collections.emptyMap()); + final ScanQuery query = Druids.newScanQueryBuilder() + .dataSource(new SystemTableDataSource(descriptor.getTableName())) + .eternityInterval() + .columns(List.of("datasource", "size")) + .resultFormat(ScanQuery.ResultFormat.RESULT_FORMAT_COMPACTED_LIST) + .build(); final List<ScanResultValue> results = client.createRunner(query, AUTHENTICATION_RESULT, false) .run(QueryPlus.wrap(query), ResponseContext.createEmpty()) @@ -392,10 +393,13 @@ public class SystemTableQueryClientTest Assertions.assertEquals(1, results.size()); Assertions.assertEquals(1, ((List<?>) results.get(0).getEvents()).size()); Mockito.verifyNoInteractions(nodeLocator, directClientFactory); - Mockito.verify(localQueryHandler).createRunner( + final ArgumentCaptor<ScanQuery> localQuery = ArgumentCaptor.forClass(ScanQuery.class); + Mockito.verify(localQueryHandler).resolveDataSource(localQuery.capture(), Mockito.same(AUTHENTICATION_RESULT)); + Assertions.assertEquals(List.of("datasource", "size"), localQuery.getValue().getColumns()); + Mockito.verify(localQueryHandler, Mockito.never()).createRunner( ArgumentMatchers.any(), - Mockito.same(escalatedAuthenticationResult), - Mockito.eq(true) + ArgumentMatchers.any(), + ArgumentMatchers.anyBoolean() ); } diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProvider.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProvider.java index 214ddef3d69..8aaaa0c6ff4 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProvider.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProvider.java @@ -107,14 +107,29 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider ) { return Iterables.transform( - Iterables.filter( - getRawRows(segmentMetadataCacheProvider.get(), metadataView, getDataSourceFilter(filters)), - Objects::nonNull - ), + getRawRows(filters, internalAuthenticationResult), row -> projectRow(row, null, jsonMapper) ); } + @Override + public Iterable<Object[]> getRawRows( + final List<DimFilter> filters, + final AuthenticationResult internalAuthenticationResult + ) + { + return Iterables.filter( + getRawRows(segmentMetadataCacheProvider.get(), metadataView, getDataSourceFilter(filters)), + Objects::nonNull + ); + } + + @Override + public Object[] projectRow(final Object[] row, @Nullable final int[] projects) + { + return projectRow(row, projects, jsonMapper); + } + /** Returns the unprojected rows shared by the native provider and the Bindable system-table implementation. */ static Iterable<Object[]> getRawRows( final BrokerSegmentMetadataCache segmentMetadataCache, diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProviderTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProviderTest.java index 8d8f54fa01f..17c8c792f88 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProviderTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProviderTest.java @@ -87,6 +87,28 @@ public class SegmentsTableDataProviderTest Assertions.assertEquals("[\"metric1\"]", row[16]); } + @Test + public void testProjectionSkipsUnrequestedJsonColumns() throws Exception + { + final ObjectMapper mapper = Mockito.mock(ObjectMapper.class); + final SegmentsTableDataProvider provider = new SegmentsTableDataProvider( + () -> Mockito.mock(BrokerSegmentMetadataCache.class), + Mockito.mock(MetadataSegmentView.class), + mapper + ); + final Object[] rawRow = new Object[SegmentsTableDescriptor.ROW_SIGNATURE.size()]; + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("datasource")] = "foo"; + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("dimensions")] = List.of("unused"); + + final Object[] projectedRow = provider.projectRow( + rawRow, + new int[]{SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("datasource")} + ); + + Assertions.assertArrayEquals(new Object[]{"foo"}, projectedRow); + Mockito.verify(mapper, Mockito.never()).writeValueAsString(Mockito.any()); + } + private static List<Object[]> toRows(final Iterable<Object[]> rows) { final List<Object[]> result = new java.util.ArrayList<>(); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
