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 2f7662c7903ce0f5fd43f37ac0ce8b7914938c09 Author: Frank Chen <[email protected]> AuthorDate: Fri Sep 4 14:07:19 2026 +0800 perf(sql): push down sys.server_segments filters --- docs/querying/sql-metadata-tables.md | 10 +-- .../query/NativeSysServerSegmentsQueryTest.java | 6 +- .../table/ServerSegmentsTableDataProvider.java | 86 +++++++++++++++++++++- .../table/ServerSegmentsTableDescriptor.java | 1 + .../ServerSegmentsTableDataProviderTest.java | 76 ++++++++++++++++++- .../druid/sql/calcite/schema/SystemSchema.java | 1 + .../druid/sql/calcite/schema/SystemSchemaTest.java | 15 ++-- 7 files changed, 177 insertions(+), 18 deletions(-) diff --git a/docs/querying/sql-metadata-tables.md b/docs/querying/sql-metadata-tables.md index c4a4d56d4c5..61c40d5dede 100644 --- a/docs/querying/sql-metadata-tables.md +++ b/docs/querying/sql-metadata-tables.md @@ -298,23 +298,23 @@ SELECT * FROM sys.servers; ### SERVER_SEGMENTS table -SERVER_SEGMENTS is used to join servers with segments table +SERVER_SEGMENTS maps servers to their loaded segments and datasources. |Column|Type|Notes| |------|-----|-----| |server|VARCHAR|Server name in format host:port (Primary key of [servers table](#servers-table))| |segment_id|VARCHAR|Segment identifier (Primary key of [segments table](#segments-table))| +|datasource|VARCHAR|Datasource name| JOIN between "servers" and "segments" can be used to query the number of segments for a specific datasource, grouped by server, example query: ```sql -SELECT count(segments.segment_id) as num_segments from sys.segments as segments -INNER JOIN sys.server_segments as server_segments -ON segments.segment_id = server_segments.segment_id +SELECT count(server_segments.segment_id) as num_segments +FROM sys.server_segments as server_segments INNER JOIN sys.servers as servers ON servers.server = server_segments.server -WHERE segments.datasource = 'wikipedia' +WHERE server_segments.datasource = 'wikipedia' GROUP BY servers.server; ``` diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerSegmentsQueryTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerSegmentsQueryTest.java index 51d8843fa76..c3b8373ff6d 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerSegmentsQueryTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysServerSegmentsQueryTest.java @@ -55,12 +55,12 @@ public class NativeSysServerSegmentsQueryTest extends EmbeddedClusterTestBase public void testNativeAggregationUsesBrokerLocalProvider(final String plannerStrategy) { final String result = cluster.runSql( - "SELECT COUNT(*), COUNT(DISTINCT server), COUNT(DISTINCT segment_id) " - + "FROM sys.server_segments", + "SELECT COUNT(*), COUNT(DISTINCT server), COUNT(DISTINCT segment_id), COUNT(DISTINCT datasource) " + + "FROM sys.server_segments WHERE datasource = 'wikipedia'", nativeQueryContext(plannerStrategy) ); - Assertions.assertEquals("0,0,0", result); + Assertions.assertEquals("0,0,0,0", result); } private static Map<String, Object> nativeQueryContext(final String plannerStrategy) diff --git a/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDataProvider.java b/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDataProvider.java index 6150c395b8c..c4fe8d2019a 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDataProvider.java +++ b/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDataProvider.java @@ -20,24 +20,39 @@ package org.apache.druid.server.system.table; import com.google.common.collect.Iterables; +import com.google.common.collect.Sets; import com.google.inject.Inject; +import org.apache.druid.client.ImmutableDruidDataSource; import org.apache.druid.client.ImmutableDruidServer; import org.apache.druid.client.TimelineServerView; import org.apache.druid.query.BatchedInlineDataSource; import org.apache.druid.query.DataSource; import org.apache.druid.query.filter.DimFilter; +import org.apache.druid.query.filter.EqualityFilter; +import org.apache.druid.query.filter.InDimFilter; +import org.apache.druid.query.filter.OrDimFilter; +import org.apache.druid.query.filter.SelectorDimFilter; +import org.apache.druid.query.filter.TypedInFilter; import org.apache.druid.server.security.AuthenticationResult; import org.apache.druid.server.security.AuthorizationUtils; import org.apache.druid.server.security.AuthorizerMapper; import org.apache.druid.timeline.DataSegment; +import javax.annotation.Nullable; import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.Set; /** Native row supplier for {@code sys.server_segments}. */ public class ServerSegmentsTableDataProvider implements SystemTableDataProvider { + private static final List<SystemTablePushdownFilter> PUSHDOWN_FILTERS = List.of( + new SystemTablePushdownFilter("server", null), + new SystemTablePushdownFilter("segment_id", null), + new SystemTablePushdownFilter("datasource", null) + ); + private final TimelineServerView serverView; private final AuthorizerMapper authorizerMapper; @@ -51,6 +66,12 @@ public class ServerSegmentsTableDataProvider implements SystemTableDataProvider this.authorizerMapper = authorizerMapper; } + @Override + public List<SystemTablePushdownFilter> getPushdownFilters() + { + return PUSHDOWN_FILTERS; + } + @Override public Optional<DataSource> getAuthorizedDataSource( final SystemTableQueryRequest request, @@ -75,21 +96,49 @@ public class ServerSegmentsTableDataProvider implements SystemTableDataProvider final AuthenticationResult authenticationResult ) { + final Set<String> serverFilter = getStringFilter(filters, "server"); + final Set<String> segmentFilter = getStringFilter(filters, "segment_id"); + final Set<String> dataSourceFilter = getStringFilter(filters, "datasource"); + final Iterable<ImmutableDruidServer> servers = serverFilter == null + ? serverView.getDruidServers() + : Iterables.filter( + serverView.getDruidServers(), + server -> serverFilter.contains(server.getHost()) + ); final Iterable<Iterable<Object[]>> rowsByServer = Iterables.transform( - serverView.getDruidServers(), - server -> getAuthorizedRows(server, authenticationResult) + servers, + server -> getAuthorizedRows(server, segmentFilter, dataSourceFilter, authenticationResult) ); return Iterables.concat(rowsByServer); } private Iterable<Object[]> getAuthorizedRows( final ImmutableDruidServer server, + @Nullable final Set<String> segmentFilter, + @Nullable final Set<String> dataSourceFilter, final AuthenticationResult authenticationResult ) { + final Iterable<ImmutableDruidDataSource> dataSources = dataSourceFilter == null + ? server.getDataSources() + : Iterables.filter( + server.getDataSources(), + dataSource -> dataSourceFilter.contains( + dataSource.getName() + ) + ); + final Iterable<DataSegment> segments = Iterables.concat( + Iterables.transform(dataSources, ImmutableDruidDataSource::getSegments) + ); + final Iterable<DataSegment> filteredSegments = segmentFilter == null + ? segments + : Iterables.filter( + segments, + segment -> segmentFilter.contains(segment.getId().toString()) + ); final Iterable<DataSegment> authorizedSegments = AuthorizationUtils.filterAuthorizedResources( authenticationResult, - server.iterateAllSegments(), + filteredSegments, segment -> Collections.singletonList( AuthorizationUtils.DATASOURCE_READ_RA_GENERATOR.apply(segment.getDataSource()) ), @@ -97,7 +146,36 @@ public class ServerSegmentsTableDataProvider implements SystemTableDataProvider ); return Iterables.transform( authorizedSegments, - segment -> new Object[]{server.getHost(), segment.getId().toString()} + segment -> new Object[]{server.getHost(), segment.getId().toString(), segment.getDataSource()} ); } + + @Nullable + private static Set<String> getStringFilter(final List<DimFilter> filters, final String column) + { + Set<String> result = null; + for (final DimFilter filter : filters) { + if (isStringValuesFilter(filter) + && column.equals(SystemTablePushdownFilter.getStringValuesColumn(filter))) { + final Set<String> values = SystemTablePushdownFilter.getStringValues(filter); + result = result == null ? values : Sets.intersection(result, values).immutableCopy(); + } + } + return result; + } + + private static boolean isStringValuesFilter(final DimFilter filter) + { + if (filter instanceof SelectorDimFilter + || filter instanceof EqualityFilter + || filter instanceof InDimFilter + || filter instanceof TypedInFilter) { + return true; + } + if (filter instanceof OrDimFilter or) { + return !or.getFields().isEmpty() + && or.getFields().stream().allMatch(ServerSegmentsTableDataProvider::isStringValuesFilter); + } + return false; + } } diff --git a/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDescriptor.java b/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDescriptor.java index e973590d494..f8580af18c9 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDescriptor.java +++ b/server/src/main/java/org/apache/druid/server/system/table/ServerSegmentsTableDescriptor.java @@ -40,6 +40,7 @@ public class ServerSegmentsTableDescriptor implements SystemTableDescriptor .builder() .add("server", ColumnType.STRING) .add("segment_id", ColumnType.STRING) + .add("datasource", ColumnType.STRING) .build(); private static final Set<NodeRole> NODE_ROLES = Set.of(NodeRole.BROKER); diff --git a/server/src/test/java/org/apache/druid/server/system/ServerSegmentsTableDataProviderTest.java b/server/src/test/java/org/apache/druid/server/system/ServerSegmentsTableDataProviderTest.java index c8fd6df6dcc..bde04ace487 100644 --- a/server/src/test/java/org/apache/druid/server/system/ServerSegmentsTableDataProviderTest.java +++ b/server/src/test/java/org/apache/druid/server/system/ServerSegmentsTableDataProviderTest.java @@ -25,6 +25,7 @@ import org.apache.druid.discovery.NodeRole; import org.apache.druid.java.util.common.Intervals; import org.apache.druid.query.BatchedInlineDataSource; import org.apache.druid.query.DataSource; +import org.apache.druid.query.filter.SelectorDimFilter; import org.apache.druid.segment.column.RowSignature; import org.apache.druid.server.coordination.ServerType; import org.apache.druid.server.security.Access; @@ -98,6 +99,46 @@ public class ServerSegmentsTableDataProviderTest EasyMock.verify(serverView); } + @Test + public void testPushesDownServerSegmentAndDatasourceFilters() + { + final DataSegment selectedSegment = segment("selected", "2024-01-01/2024-01-02"); + final DataSegment otherSegment = segment("other", "2024-01-02/2024-01-03"); + final DruidServer selectedServer = server("selected:8083", selectedSegment, otherSegment); + final DruidServer otherServer = server("other:8083", selectedSegment); + final TimelineServerView serverView = EasyMock.mock(TimelineServerView.class); + EasyMock.expect(serverView.getDruidServers()).andReturn( + List.of(selectedServer.toImmutableDruidServer(), otherServer.toImmutableDruidServer()) + ).once(); + EasyMock.replay(serverView); + + final ServerSegmentsTableDataProvider provider = new ServerSegmentsTableDataProvider( + serverView, + allowAllAuthorizerMapper() + ); + final List<Object[]> rows = toRows( + provider.getRows( + List.of( + new SelectorDimFilter("server", "selected:8083", null), + new SelectorDimFilter("segment_id", selectedSegment.getId().toString(), null), + new SelectorDimFilter("datasource", "selected", null) + ), + AUTHENTICATION_RESULT + ) + ); + + Assertions.assertEquals( + List.of("server", "segment_id", "datasource"), + provider.getPushdownFilters().stream().map(filter -> filter.key()).toList() + ); + Assertions.assertEquals(1, rows.size()); + Assertions.assertArrayEquals( + new Object[]{"selected:8083", selectedSegment.getId().toString(), "selected"}, + rows.get(0) + ); + EasyMock.verify(serverView); + } + @Test public void testDescriptorRejectsRequestWithoutStateRead() { @@ -120,7 +161,10 @@ public class ServerSegmentsTableDataProviderTest Assertions.assertEquals(Set.of(NodeRole.BROKER), descriptor.getNodeRoles()); Assertions.assertEquals(SystemTableRoutingMode.LOCAL_ONLY, descriptor.getRoutingMode()); - Assertions.assertEquals(List.of("server", "segment_id"), descriptor.getRowSignature().getColumnNames()); + Assertions.assertEquals( + List.of("server", "segment_id", "datasource"), + descriptor.getRowSignature().getColumnNames() + ); } private static AuthorizerMapper authorizerMapper(final boolean allowState) @@ -138,6 +182,36 @@ public class ServerSegmentsTableDataProviderTest }; } + private static AuthorizerMapper allowAllAuthorizerMapper() + { + return new AuthorizerMapper(null) + { + @Override + public org.apache.druid.server.security.Authorizer getAuthorizer(final String name) + { + return (authenticationResult, resource, action) -> Access.OK; + } + }; + } + + private static DruidServer server(final String host, final DataSegment... segments) + { + final DruidServer server = new DruidServer( + host, + host, + null, + 1_000L, + null, + ServerType.HISTORICAL, + "default", + 0 + ); + for (final DataSegment segment : segments) { + server.addDataSegment(segment); + } + return server; + } + private static DataSegment segment(final String dataSource, final String interval) { return DataSegment.builder() diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java index b43e0a415dc..fe716c1d6c8 100644 --- a/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SystemSchema.java @@ -719,6 +719,7 @@ public class SystemSchema extends AbstractTableSchema Object[] row = new Object[serverSegmentsTableSize]; row[0] = druidServer.getHost(); row[1] = segment.getId().toString(); + row[2] = segment.getDataSource(); rows.add(row); } } diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java index f63d04b71f4..445cb2ea169 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemSchemaTest.java @@ -1360,11 +1360,11 @@ public class SystemSchemaTest extends CalciteTestBase //server_segments table is the join of servers and segments table // it will have 5 rows as follows - // localhost:0000 | test1_2010-01-01T00:00:00.000Z_2011-01-01T00:00:00.000Z_version1(segment1) - // localhost:0000 | test2_2011-01-01T00:00:00.000Z_2012-01-01T00:00:00.000Z_version2(segment2) - // server2:1234 | test3_2012-01-01T00:00:00.000Z_2013-01-01T00:00:00.000Z_version3(segment3) - // server2:1234 | test4_2017-01-01T00:00:00.000Z_2018-01-01T00:00:00.000Z_version4(segment4) - // server2:1234 | test5_2017-01-01T00:00:00.000Z_2018-01-01T00:00:00.000Z_version5(segment5) + // localhost:0000 | test1_2010-01-01T00:00:00.000Z_2011-01-01T00:00:00.000Z_version1 | test1 + // localhost:0000 | test2_2011-01-01T00:00:00.000Z_2012-01-01T00:00:00.000Z_version2 | test2 + // server2:1234 | test3_2012-01-01T00:00:00.000Z_2013-01-01T00:00:00.000Z_version3 | test3 + // server2:1234 | test4_2017-01-01T00:00:00.000Z_2018-01-01T00:00:00.000Z_version4 | test4 + // server2:1234 | test5_2017-01-01T00:00:00.000Z_2018-01-01T00:00:00.000Z_version5 | test5 final List<Object[]> rows = serverSegmentsTable.scan(dataContext).toList(); Assertions.assertEquals(5, rows.size()); @@ -1372,22 +1372,27 @@ public class SystemSchemaTest extends CalciteTestBase Object[] row0 = rows.get(0); Assertions.assertEquals("localhost:0000", row0[0]); Assertions.assertEquals("test1_2010-01-01T00:00:00.000Z_2011-01-01T00:00:00.000Z_version1", row0[1].toString()); + Assertions.assertEquals("test1", row0[2]); Object[] row1 = rows.get(1); Assertions.assertEquals("localhost:0000", row1[0]); Assertions.assertEquals("test2_2011-01-01T00:00:00.000Z_2012-01-01T00:00:00.000Z_version2", row1[1].toString()); + Assertions.assertEquals("test2", row1[2]); Object[] row2 = rows.get(2); Assertions.assertEquals("server2:1234", row2[0]); Assertions.assertEquals("test3_2012-01-01T00:00:00.000Z_2013-01-01T00:00:00.000Z_version3_2", row2[1].toString()); + Assertions.assertEquals("test3", row2[2]); Object[] row3 = rows.get(3); Assertions.assertEquals("server2:1234", row3[0]); Assertions.assertEquals("test4_2014-01-01T00:00:00.000Z_2015-01-01T00:00:00.000Z_version4", row3[1].toString()); + Assertions.assertEquals("test4", row3[2]); Object[] row4 = rows.get(4); Assertions.assertEquals("server2:1234", row4[0]); Assertions.assertEquals("test5_2015-01-01T00:00:00.000Z_2016-01-01T00:00:00.000Z_version5", row4[1].toString()); + Assertions.assertEquals("test5", row4[2]); // Verify value types. verifyTypes(rows, ServerSegmentsTableDescriptor.ROW_SIGNATURE); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
