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 559c3d1857e04d6f68ef8c21d2f88c4bdd1cac33 Author: Frank Chen <[email protected]> AuthorDate: Mon Aug 31 15:48:21 2026 +0800 feat(sql): add native query support for sys.segments --- docs/querying/sql-metadata-tables.md | 1 + .../embedded/query/NativeSysSegmentsQueryTest.java | 85 ++++++ .../system/handler/SystemTableQueryClient.java | 29 ++ .../server/system/module/SystemTableModule.java | 4 + .../system/table/SegmentsTableDescriptor.java | 111 +++++++ .../system/table/SystemTableRoutingMode.java | 3 + .../system/handler/SystemTableQueryClientTest.java | 59 ++++ .../system/table/SegmentsTableDescriptorTest.java | 102 +++++++ .../main/java/org/apache/druid/cli/CliBroker.java | 9 + .../test/java/org/apache/druid/cli/MainTest.java | 16 + .../sql/calcite/schema/NativeSegmentsTable.java | 70 +++++ .../calcite/schema/SegmentsTableDataProvider.java | 279 +++++++++++++++++ .../druid/sql/calcite/schema/SystemSchema.java | 330 ++------------------- .../schema/SegmentsTableDataProviderTest.java | 96 ++++++ .../schema/SystemTableDataProviderTest.java | 6 + 15 files changed, 892 insertions(+), 308 deletions(-) diff --git a/docs/querying/sql-metadata-tables.md b/docs/querying/sql-metadata-tables.md index 8b47335617f..c4a4d56d4c5 100644 --- a/docs/querying/sql-metadata-tables.md +++ b/docs/querying/sql-metadata-tables.md @@ -168,6 +168,7 @@ execution: |Table|Source of rows| |-----|--------------| +|[`sys.segments`](#segments-table)|The Broker that receives the SQL query. It combines that Broker's segment metadata cache with its Coordinator metadata view, and exact `datasource` equality/`IN` filters are pushed into the local scan. Native execution does not fan out to other Brokers.| |[`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, diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysSegmentsQueryTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysSegmentsQueryTest.java new file mode 100644 index 00000000000..986509e43ed --- /dev/null +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/query/NativeSysSegmentsQueryTest.java @@ -0,0 +1,85 @@ +/* + * 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.junit.jupiter.api.Assertions; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.util.Map; + +public class NativeSysSegmentsQueryTest extends QueryTestBase +{ + private String dataSource; + + @Override + public void beforeAll() + { + dataSource = ingestBasicData(); + } + + @ParameterizedTest(name = "plannerStrategy = {0}") + @ValueSource(strings = { + QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_COUPLED, + QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_DECOUPLED + }) + public void testNativeSegmentAggregations(final String plannerStrategy) + { + final String result = cluster.runSql( + "SELECT COUNT(*), COUNT(DISTINCT segment_id), COUNT(DISTINCT datasource), SUM(1) " + + "FROM sys.segments WHERE datasource = '" + dataSource + "'", + nativeQueryContext(plannerStrategy) + ); + + Assertions.assertEquals("10,10,1,10", result); + } + + @ParameterizedTest(name = "plannerStrategy = {0}") + @ValueSource(strings = { + QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_COUPLED, + QueryContexts.NATIVE_QUERY_SQL_PLANNING_MODE_DECOUPLED + }) + public void testNativeNestedSegmentAggregation(final String plannerStrategy) + { + final String result = cluster.runSql( + "SELECT COUNT(*) FROM (" + + "SELECT segment_id FROM sys.segments WHERE datasource = '" + dataSource + "' GROUP BY segment_id" + + ")", + nativeQueryContext(plannerStrategy) + ); + + Assertions.assertEquals("10", result); + } + + private static Map<String, Object> 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/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java index 7330e51553e..58577802cbb 100644 --- a/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java +++ b/server/src/main/java/org/apache/druid/server/system/handler/SystemTableQueryClient.java @@ -282,6 +282,11 @@ 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( @@ -305,6 +310,30 @@ 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, diff --git a/server/src/main/java/org/apache/druid/server/system/module/SystemTableModule.java b/server/src/main/java/org/apache/druid/server/system/module/SystemTableModule.java index 6e25c527f3e..c31cf51da86 100644 --- a/server/src/main/java/org/apache/druid/server/system/module/SystemTableModule.java +++ b/server/src/main/java/org/apache/druid/server/system/module/SystemTableModule.java @@ -25,6 +25,7 @@ import com.google.inject.multibindings.MapBinder; import org.apache.druid.guice.DruidBinders; import org.apache.druid.guice.LazySingleton; import org.apache.druid.query.SystemTableDataSource; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; import org.apache.druid.server.system.table.ServerPropertiesTableDataProvider; import org.apache.druid.server.system.table.ServerPropertiesTableDescriptor; import org.apache.druid.server.system.table.SystemTableDataProvider; @@ -48,6 +49,9 @@ public class SystemTableModule implements Module final MapBinder<String, SystemTableDescriptor> descriptorBinder = MapBinder.newMapBinder(binder, String.class, SystemTableDescriptor.class); descriptorBinder.addBinding(ServerPropertiesTableDescriptor.TABLE_NAME) .toInstance(new ServerPropertiesTableDescriptor()); + descriptorBinder.addBinding(SegmentsTableDescriptor.TABLE_NAME) + .toInstance(new SegmentsTableDescriptor()); + final MapBinder<String, SystemTableDataProvider> dataProviderBinder = MapBinder.newMapBinder(binder, String.class, SystemTableDataProvider.class); dataProviderBinder.addBinding(ServerPropertiesTableDescriptor.TABLE_NAME) .to(ServerPropertiesTableDataProvider.class) diff --git a/server/src/main/java/org/apache/druid/server/system/table/SegmentsTableDescriptor.java b/server/src/main/java/org/apache/druid/server/system/table/SegmentsTableDescriptor.java new file mode 100644 index 00000000000..80822fbb605 --- /dev/null +++ b/server/src/main/java/org/apache/druid/server/system/table/SegmentsTableDescriptor.java @@ -0,0 +1,111 @@ +/* + * 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.table; + +import org.apache.druid.discovery.NodeRole; +import org.apache.druid.segment.column.ColumnType; +import org.apache.druid.segment.column.RowSignature; +import org.apache.druid.server.security.AuthenticationResult; +import org.apache.druid.server.security.AuthorizationUtils; +import org.apache.druid.server.security.AuthorizerMapper; + +import java.util.Collections; +import java.util.Set; + +/** Descriptor for the native {@code sys.segments} table. */ +public class SegmentsTableDescriptor implements SystemTableDescriptor +{ + public static final String TABLE_NAME = "segments"; + public static final RowSignature ROW_SIGNATURE = RowSignature + .builder() + .add("segment_id", ColumnType.STRING) + .add("datasource", ColumnType.STRING) + .add("start", ColumnType.STRING) + .add("end", ColumnType.STRING) + .add("size", ColumnType.LONG) + .add("version", ColumnType.STRING) + .add("partition_num", ColumnType.LONG) + .add("num_replicas", ColumnType.LONG) + .add("num_rows", ColumnType.LONG) + .add("is_active", ColumnType.LONG) + .add("is_published", ColumnType.LONG) + .add("is_available", ColumnType.LONG) + .add("is_realtime", ColumnType.LONG) + .add("is_overshadowed", ColumnType.LONG) + .add("shard_spec", ColumnType.STRING) + .add("dimensions", ColumnType.STRING) + .add("metrics", ColumnType.STRING) + .add("projections", ColumnType.STRING) + .add("last_compaction_state", ColumnType.STRING) + .add("replication_factor", ColumnType.LONG) + .build(); + + private static final Set<NodeRole> NODE_ROLES = Set.of(NodeRole.BROKER); + private static final int DATASOURCE_COLUMN = ROW_SIGNATURE.indexOf("datasource"); + private static final SystemTableRowAuthorizer ROW_AUTHORIZER = new SystemTableRowAuthorizer() + { + @Override + public Iterable<Object[]> filterAuthorizedRows( + final Iterable<Object[]> rows, + final AuthenticationResult authenticationResult, + final AuthorizerMapper authorizerMapper + ) + { + return AuthorizationUtils.filterAuthorizedResources( + authenticationResult, + rows, + row -> Collections.singletonList( + AuthorizationUtils.DATASOURCE_READ_RA_GENERATOR.apply((String) row[DATASOURCE_COLUMN]) + ), + authorizerMapper + ); + } + }; + + @Override + public String getTableName() + { + return TABLE_NAME; + } + + @Override + public Set<NodeRole> getNodeRoles() + { + return NODE_ROLES; + } + + @Override + public SystemTableRoutingMode getRoutingMode() + { + return SystemTableRoutingMode.LOCAL_ONLY; + } + + @Override + public RowSignature getRowSignature() + { + return ROW_SIGNATURE; + } + + @Override + public SystemTableRowAuthorizer getRowAuthorizer() + { + return ROW_AUTHORIZER; + } +} diff --git a/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java b/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java index fa6c1c69ba6..a8589501c74 100644 --- a/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java +++ b/server/src/main/java/org/apache/druid/server/system/table/SystemTableRoutingMode.java @@ -22,6 +22,9 @@ package org.apache.druid.server.system.table; /** Defines how the Broker selects nodes that contribute rows to a system table. */ public enum SystemTableRoutingMode { + /** Only the node receiving the query contributes its local rows. */ + LOCAL_ONLY, + /** Every discovered node for the descriptor's roles contributes an independent set of rows. */ ALL_NODES, diff --git a/server/src/test/java/org/apache/druid/server/system/handler/SystemTableQueryClientTest.java b/server/src/test/java/org/apache/druid/server/system/handler/SystemTableQueryClientTest.java index 3e7273efc20..45e3bfafae8 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 @@ -81,6 +81,7 @@ import org.apache.druid.server.security.Escalator; import org.apache.druid.server.security.ForbiddenException; import org.apache.druid.server.security.NoopEscalator; import org.apache.druid.server.system.SystemTableNotLeaderException; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; import org.apache.druid.server.system.table.ServerPropertiesTableDescriptor; import org.apache.druid.server.system.table.SystemTableDescriptor; import org.apache.druid.server.system.table.SystemTableRoutingMode; @@ -340,6 +341,64 @@ public class SystemTableQueryClientTest ); } + /** A local-only system table bypasses discovery and remote clients entirely. */ + @Test + public void testLocalOnlySystemTableExecutesOnReceivingBroker() + { + final SystemTableDescriptor descriptor = new SegmentsTableDescriptor(); + 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) + ); + + final QuerySegmentWalker querySegmentWalker = Mockito.mock(QuerySegmentWalker.class); + Mockito.when(querySegmentWalker.getQueryRunnerForIntervals(ArgumentMatchers.any(), ArgumentMatchers.any())) + .thenAnswer(ignored -> passthroughRunner(descriptor.getRowSignature())); + final SystemTableQueryClient client = new SystemTableQueryClient( + nodeLocator, + directClientFactory, + Mockito.mock(QueryScheduler.class), + querySegmentWalker, + Map.of(descriptor.getTableName(), descriptor), + new AuthorizerMapper(Map.of("allow", new AllowAllAuthorizer(null))), + localQueryHandler, + escalator, + nonMatchingSelfNode() + ); + final ScanQuery query = query(descriptor, Collections.emptyMap()); + + final List<ScanResultValue> results = client.createRunner(query, AUTHENTICATION_RESULT, false) + .run(QueryPlus.wrap(query), ResponseContext.createEmpty()) + .toList(); + + Assertions.assertEquals(1, results.size()); + Assertions.assertEquals(1, ((List<?>) results.get(0).getEvents()).size()); + Mockito.verifyNoInteractions(nodeLocator, directClientFactory); + Mockito.verify(localQueryHandler).createRunner( + ArgumentMatchers.any(), + Mockito.same(escalatedAuthenticationResult), + Mockito.eq(true) + ); + } + /** Delayed node responses verify concurrent fanout through the real {@link DirectDruidClient} transport path. */ @Test public void testDirectClientsStartRequestsBeforeWaitingForDelayedResponses() throws Exception diff --git a/server/src/test/java/org/apache/druid/server/system/table/SegmentsTableDescriptorTest.java b/server/src/test/java/org/apache/druid/server/system/table/SegmentsTableDescriptorTest.java new file mode 100644 index 00000000000..a317c190317 --- /dev/null +++ b/server/src/test/java/org/apache/druid/server/system/table/SegmentsTableDescriptorTest.java @@ -0,0 +1,102 @@ +/* + * 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.table; + +import org.apache.druid.discovery.NodeRole; +import org.apache.druid.java.util.common.Intervals; +import org.apache.druid.server.security.Access; +import org.apache.druid.server.security.AuthConfig; +import org.apache.druid.server.security.AuthenticationResult; +import org.apache.druid.server.security.Authorizer; +import org.apache.druid.server.security.AuthorizerMapper; +import org.apache.druid.timeline.DataSegment; +import org.apache.druid.timeline.SegmentId; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Set; + +public class SegmentsTableDescriptorTest +{ + private static final AuthenticationResult AUTHENTICATION_RESULT = + new AuthenticationResult("test-user", AuthConfig.ALLOW_ALL_NAME, null, null); + + @Test + public void testDescriptorDefinesBrokerLocalSegmentsTable() + { + final SegmentsTableDescriptor descriptor = new SegmentsTableDescriptor(); + + Assertions.assertEquals("segments", descriptor.getTableName()); + Assertions.assertEquals(Set.of(NodeRole.BROKER), descriptor.getNodeRoles()); + Assertions.assertEquals(SystemTableRoutingMode.LOCAL_ONLY, descriptor.getRoutingMode()); + Assertions.assertEquals(20, descriptor.getRowSignature().size()); + Assertions.assertEquals("datasource", descriptor.getRowSignature().getColumnNames().get(1)); + } + + @Test + public void testDescriptorFiltersRowsByDatasourceReadPermission() + { + final DataSegment allowedSegment = DataSegment.builder( + SegmentId.of("allowed", Intervals.of("2000/2001"), "v", null) + ).build(); + final DataSegment deniedSegment = DataSegment.builder( + SegmentId.of("denied", Intervals.of("2000/2001"), "v", null) + ).build(); + final List<Object[]> rows = List.of( + row(allowedSegment.getDataSource()), + row(deniedSegment.getDataSource()) + ); + final AuthorizerMapper authorizerMapper = new AuthorizerMapper(null) + { + @Override + public Authorizer getAuthorizer(final String name) + { + return (authenticationResult, resource, action) -> + "allowed".equals(resource.getName()) ? Access.OK : Access.DENIED; + } + }; + + final List<Object[]> authorizedRows = toList( + new SegmentsTableDescriptor().getRowAuthorizer().filterAuthorizedRows( + rows, + AUTHENTICATION_RESULT, + authorizerMapper + ) + ); + + Assertions.assertEquals(1, authorizedRows.size()); + Assertions.assertEquals("allowed", authorizedRows.get(0)[1]); + } + + private static Object[] row(final String datasource) + { + final Object[] row = new Object[SegmentsTableDescriptor.ROW_SIGNATURE.size()]; + row[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("datasource")] = datasource; + return row; + } + + private static List<Object[]> toList(final Iterable<Object[]> rows) + { + final List<Object[]> result = new java.util.ArrayList<>(); + rows.forEach(result::add); + return result; + } +} diff --git a/services/src/main/java/org/apache/druid/cli/CliBroker.java b/services/src/main/java/org/apache/druid/cli/CliBroker.java index 2723b496b97..e37b48e5e4c 100644 --- a/services/src/main/java/org/apache/druid/cli/CliBroker.java +++ b/services/src/main/java/org/apache/druid/cli/CliBroker.java @@ -24,6 +24,7 @@ import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableSet; import com.google.inject.Key; import com.google.inject.Module; +import com.google.inject.multibindings.MapBinder; import com.google.inject.name.Names; import org.apache.druid.client.BrokerSegmentWatcherConfig; import org.apache.druid.client.BrokerServerView; @@ -81,7 +82,10 @@ import org.apache.druid.server.http.SelfDiscoveryResource; import org.apache.druid.server.initialization.jetty.JettyServerInitializer; import org.apache.druid.server.metrics.SubqueryCountStatsProvider; import org.apache.druid.server.router.TieredBrokerConfig; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; +import org.apache.druid.server.system.table.SystemTableDataProvider; import org.apache.druid.sql.calcite.schema.MetadataSegmentView; +import org.apache.druid.sql.calcite.schema.SegmentsTableDataProvider; import org.apache.druid.sql.guice.SqlModule; import org.apache.druid.storage.local.LocalTmpStorageConfig; import org.apache.druid.timeline.PruneLoadSpec; @@ -126,6 +130,11 @@ public class CliBroker extends ServerRunnable binder -> { validateCentralizedDatasourceSchemaConfig(getProperties()); + MapBinder.newMapBinder(binder, String.class, SystemTableDataProvider.class) + .addBinding(SegmentsTableDescriptor.TABLE_NAME) + .to(SegmentsTableDataProvider.class) + .in(LazySingleton.class); + binder.bindConstant().annotatedWith(Names.named("serviceName")).to( TieredBrokerConfig.DEFAULT_BROKER_SERVICE_NAME ); diff --git a/services/src/test/java/org/apache/druid/cli/MainTest.java b/services/src/test/java/org/apache/druid/cli/MainTest.java index 4ac8f2a035a..5c0b67334b5 100644 --- a/services/src/test/java/org/apache/druid/cli/MainTest.java +++ b/services/src/test/java/org/apache/druid/cli/MainTest.java @@ -24,8 +24,10 @@ import com.google.inject.Key; import com.google.inject.TypeLiteral; import org.apache.druid.guice.GuiceInjectors; import org.apache.druid.server.system.handler.SystemTableQueryResource; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; import org.apache.druid.server.system.table.ServerPropertiesTableDescriptor; import org.apache.druid.server.system.table.SystemTableDataProvider; +import org.apache.druid.sql.calcite.schema.SegmentsTableDataProvider; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; @@ -77,6 +79,20 @@ public class MainTest ); } + @Test + public void testBrokerNativeSegmentsProviderInjection() + { + final CliBroker broker = new CliBroker(); + final Injector injector = GuiceInjectors.makeStartupInjector(); + injector.injectMembers(broker); + + final Injector brokerInjector = broker.makeInjector(broker.getNodeRoles(new Properties())); + final SystemTableDataProvider provider = systemTableDataProviders(brokerInjector) + .get(SegmentsTableDescriptor.TABLE_NAME); + + Assertions.assertInstanceOf(SegmentsTableDataProvider.class, provider); + } + private static Map<String, SystemTableDataProvider> systemTableDataProviders(final Injector injector) { return injector.getInstance(Key.get(new TypeLiteral<>() {})); diff --git a/sql/src/main/java/org/apache/druid/sql/calcite/schema/NativeSegmentsTable.java b/sql/src/main/java/org/apache/druid/sql/calcite/schema/NativeSegmentsTable.java new file mode 100644 index 00000000000..1471a2e1eb8 --- /dev/null +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/NativeSegmentsTable.java @@ -0,0 +1,70 @@ +/* + * 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.sql.calcite.schema; + +import org.apache.calcite.plan.RelOptTable; +import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.logical.LogicalTableScan; +import org.apache.calcite.schema.Schema; +import org.apache.druid.query.DataSource; +import org.apache.druid.query.SystemTableDataSource; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; +import org.apache.druid.sql.calcite.table.DruidTable; + +/** Native-query representation of {@code sys.segments}. */ +class NativeSegmentsTable extends DruidTable +{ + private static final DataSource DATA_SOURCE = new SystemTableDataSource(SegmentsTableDescriptor.TABLE_NAME); + + NativeSegmentsTable() + { + super(SegmentsTableDescriptor.ROW_SIGNATURE); + } + + @Override + public DataSource getDataSource() + { + return DATA_SOURCE; + } + + @Override + public boolean isJoinable() + { + return false; + } + + @Override + public boolean isBroadcast() + { + return false; + } + + @Override + public Schema.TableType getJdbcTableType() + { + return Schema.TableType.SYSTEM_TABLE; + } + + @Override + public RelNode toRel(final RelOptTable.ToRelContext context, final RelOptTable table) + { + return LogicalTableScan.create(context.getCluster(), table, context.getTableHints()); + } +} 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 new file mode 100644 index 00000000000..214ddef3d69 --- /dev/null +++ b/sql/src/main/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProvider.java @@ -0,0 +1,279 @@ +/* + * 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.sql.calcite.schema; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.FluentIterable; +import com.google.common.collect.Iterables; +import com.google.common.collect.Sets; +import com.google.inject.Inject; +import com.google.inject.Provider; +import it.unimi.dsi.fastutil.ints.IntOpenHashSet; +import it.unimi.dsi.fastutil.ints.IntSet; +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.segment.column.ValueType; +import org.apache.druid.segment.metadata.AvailableSegmentMetadata; +import org.apache.druid.server.security.AuthenticationResult; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; +import org.apache.druid.server.system.table.SystemTableDataProvider; +import org.apache.druid.server.system.table.SystemTablePushdownFilter; +import org.apache.druid.timeline.DataSegment; +import org.apache.druid.timeline.SegmentId; +import org.apache.druid.timeline.SegmentStatusInCluster; + +import javax.annotation.Nullable; +import java.util.HashSet; +import java.util.Iterator; +import java.util.List; +import java.util.Objects; +import java.util.Set; +import java.util.stream.IntStream; + +/** Native row supplier for {@code sys.segments}. */ +public class SegmentsTableDataProvider implements SystemTableDataProvider +{ + private static final long REPLICATION_FACTOR_UNKNOWN = -1L; + private static final long IS_ACTIVE_FALSE = 0L; + private static final long IS_ACTIVE_TRUE = 1L; + private static final long IS_PUBLISHED_FALSE = 0L; + private static final long IS_PUBLISHED_TRUE = 1L; + private static final long IS_AVAILABLE_TRUE = 1L; + private static final long IS_OVERSHADOWED_FALSE = 0L; + private static final long IS_OVERSHADOWED_TRUE = 1L; + + private static final int[] PROJECT_ALL = IntStream.range(0, SegmentsTableDescriptor.ROW_SIGNATURE.size()).toArray(); + private static final IntSet JSON_FIELDS = new IntOpenHashSet( + new int[]{ + SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("shard_spec"), + SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("dimensions"), + SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("metrics"), + SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("projections"), + SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("last_compaction_state") + } + ); + private static final List<SystemTablePushdownFilter> PUSHDOWN_FILTERS = List.of( + new SystemTablePushdownFilter("datasource", null) + ); + + private final Provider<BrokerSegmentMetadataCache> segmentMetadataCacheProvider; + private final MetadataSegmentView metadataView; + private final ObjectMapper jsonMapper; + + @Inject + public SegmentsTableDataProvider( + final Provider<BrokerSegmentMetadataCache> segmentMetadataCacheProvider, + final MetadataSegmentView metadataView, + final ObjectMapper jsonMapper + ) + { + this.segmentMetadataCacheProvider = segmentMetadataCacheProvider; + this.metadataView = metadataView; + this.jsonMapper = jsonMapper; + } + + @Override + public List<SystemTablePushdownFilter> getPushdownFilters() + { + return PUSHDOWN_FILTERS; + } + + @Override + public Iterable<Object[]> getRows( + final List<DimFilter> filters, + final AuthenticationResult internalAuthenticationResult + ) + { + return Iterables.transform( + Iterables.filter( + getRawRows(segmentMetadataCacheProvider.get(), metadataView, getDataSourceFilter(filters)), + Objects::nonNull + ), + row -> projectRow(row, null, jsonMapper) + ); + } + + /** Returns the unprojected rows shared by the native provider and the Bindable system-table implementation. */ + static Iterable<Object[]> getRawRows( + final BrokerSegmentMetadataCache segmentMetadataCache, + final MetadataSegmentView metadataView, + @Nullable final Set<String> dataSourceFilter + ) + { + final Set<SegmentId> segmentsAlreadySeen = dataSourceFilter == null + ? Sets.newHashSetWithExpectedSize( + segmentMetadataCache.getTotalSegments() + ) + : new HashSet<>(); + + final Iterator<SegmentStatusInCluster> metadataStoreSegments = metadataView.getSegments(dataSourceFilter); + final FluentIterable<Object[]> publishedSegments = FluentIterable + .from(() -> metadataStoreSegments) + .transform(val -> { + final DataSegment segment = val.getDataSegment(); + final AvailableSegmentMetadata availableSegmentMetadata = + segmentMetadataCache.getAvailableSegmentMetadata(segment.getDataSource(), segment.getId()); + segmentsAlreadySeen.add(segment.getId()); + + long numReplicas = 0L; + long isAvailable = 0L; + if (availableSegmentMetadata != null) { + numReplicas = availableSegmentMetadata.getNumReplicas(); + isAvailable = availableSegmentMetadata.getNumReplicas() > 0 ? IS_AVAILABLE_TRUE : IS_ACTIVE_FALSE; + } + + final long numRows; + if (segment.getTotalRows() != null) { + numRows = segment.getTotalRows().longValue(); + } else if (val.getNumRows() != null) { + numRows = val.getNumRows(); + } else if (availableSegmentMetadata != null) { + numRows = availableSegmentMetadata.getNumRows(); + } else { + numRows = 0L; + } + + final long isRealtime = val.isRealtime() ? 1 : 0; + final boolean isPublished = !val.isRealtime(); + final boolean isActive = isPublished ? !val.isOvershadowed() : val.isRealtime(); + + return new Object[]{ + segment.getId(), + segment.getDataSource(), + segment.getInterval().getStart(), + segment.getInterval().getEnd(), + segment.getSize(), + segment.getVersion(), + (long) segment.getShardSpec().getPartitionNum(), + numReplicas, + numRows, + isActive ? IS_ACTIVE_TRUE : IS_ACTIVE_FALSE, + isPublished ? IS_PUBLISHED_TRUE : IS_PUBLISHED_FALSE, + isAvailable, + isRealtime, + val.isOvershadowed() ? IS_OVERSHADOWED_TRUE : IS_OVERSHADOWED_FALSE, + segment.getShardSpec(), + segment.getDimensions(), + segment.getMetrics(), + segment.getProjections(), + segment.getLastCompactionState(), + val.getReplicationFactor() == null ? REPLICATION_FACTOR_UNKNOWN : (long) val.getReplicationFactor() + }; + }); + + final FluentIterable<Object[]> availableSegments = FluentIterable + .from(() -> segmentMetadataCache.iterateSegmentMetadata(dataSourceFilter)) + .transform(val -> { + final DataSegment segment = val.getSegment(); + if (segmentsAlreadySeen.contains(segment.getId())) { + return null; + } + return new Object[]{ + segment.getId(), + segment.getDataSource(), + segment.getInterval().getStart(), + segment.getInterval().getEnd(), + segment.getSize(), + segment.getVersion(), + (long) segment.getShardSpec().getPartitionNum(), + val.getNumReplicas(), + segment.getTotalRows() != null ? segment.getTotalRows() : val.getNumRows(), + val.isRealtime(), + IS_PUBLISHED_FALSE, + IS_AVAILABLE_TRUE, + val.isRealtime(), + IS_OVERSHADOWED_FALSE, + segment.getShardSpec(), + segment.getDimensions(), + segment.getMetrics(), + segment.getProjections(), + null, + REPLICATION_FACTOR_UNKNOWN + }; + }); + + return Iterables.unmodifiableIterable(Iterables.concat(publishedSegments, availableSegments)); + } + + /** Converts raw segment fields to the strings declared by {@link SegmentsTableDescriptor#ROW_SIGNATURE}. */ + static Object[] projectRow( + final Object[] row, + @Nullable final int[] projects, + final ObjectMapper jsonMapper + ) + { + final int[] nonNullProjects = projects == null ? PROJECT_ALL : projects; + final Object[] projectedRow = new Object[nonNullProjects.length]; + + for (int i = 0; i < nonNullProjects.length; i++) { + final int column = nonNullProjects[i]; + final Object value = row[column]; + if (SegmentsTableDescriptor.ROW_SIGNATURE.getColumnType(column).get().is(ValueType.STRING) + && value != null + && !(value instanceof String)) { + if (JSON_FIELDS.contains(column)) { + try { + projectedRow[i] = jsonMapper.writeValueAsString(value); + } + catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + } else { + projectedRow[i] = value.toString(); + } + } else { + projectedRow[i] = value; + } + } + return projectedRow; + } + + @Nullable + private static Set<String> getDataSourceFilter(final List<DimFilter> filters) + { + Set<String> result = null; + for (final DimFilter filter : filters) { + if (isStringValuesFilter(filter) && "datasource".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(SegmentsTableDataProvider::isStringValuesFilter); + } + return false; + } +} 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 94ef050892c..52379312bad 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 @@ -24,10 +24,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.base.Function; import com.google.common.collect.FluentIterable; import com.google.common.collect.Iterables; -import com.google.common.collect.Sets; import com.google.inject.Provider; -import it.unimi.dsi.fastutil.ints.IntOpenHashSet; -import it.unimi.dsi.fastutil.ints.IntSet; import org.apache.calcite.DataContext; import org.apache.calcite.linq4j.DefaultEnumerable; import org.apache.calcite.linq4j.Enumerable; @@ -62,8 +59,6 @@ import org.apache.druid.java.util.http.client.HttpClient; import org.apache.druid.rpc.indexing.OverlordClient; import org.apache.druid.segment.column.ColumnType; import org.apache.druid.segment.column.RowSignature; -import org.apache.druid.segment.column.ValueType; -import org.apache.druid.segment.metadata.AvailableSegmentMetadata; import org.apache.druid.server.DruidNode; import org.apache.druid.server.security.Action; import org.apache.druid.server.security.AuthenticationResult; @@ -74,6 +69,7 @@ import org.apache.druid.server.security.ForbiddenException; import org.apache.druid.server.security.Resource; import org.apache.druid.server.security.ResourceAction; import org.apache.druid.server.security.ResourceType; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; import org.apache.druid.sql.calcite.planner.PlannerConfig; import org.apache.druid.sql.calcite.planner.PlannerContext; import org.apache.druid.sql.calcite.run.NativeSqlEngine; @@ -84,8 +80,6 @@ import org.apache.druid.sql.http.GetQueriesResponse; import org.apache.druid.sql.http.QueryInfo; import org.apache.druid.sql.http.SqlEngineRegistry; import org.apache.druid.timeline.DataSegment; -import org.apache.druid.timeline.SegmentId; -import org.apache.druid.timeline.SegmentStatusInCluster; import javax.annotation.Nullable; import java.io.Closeable; @@ -93,7 +87,6 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; -import java.util.HashSet; import java.util.Iterator; import java.util.List; import java.util.Objects; @@ -103,39 +96,18 @@ import java.util.stream.IntStream; public class SystemSchema extends AbstractTableSchema { - public static final String SEGMENTS_TABLE = "segments"; + public static final String SEGMENTS_TABLE = SegmentsTableDescriptor.TABLE_NAME; public static final String SERVERS_TABLE = "servers"; public static final String SERVER_SEGMENTS_TABLE = "server_segments"; public static final String TASKS_TABLE = "tasks"; public static final String SUPERVISOR_TABLE = "supervisors"; public static final String QUERIES_TABLE = "queries"; - private static final Function<SegmentStatusInCluster, Iterable<ResourceAction>> - SEGMENT_STATUS_IN_CLUSTER_RA_GENERATOR = segment -> - Collections.singletonList(AuthorizationUtils.DATASOURCE_READ_RA_GENERATOR.apply( - segment.getDataSegment().getDataSource()) - ); - private static final Function<DataSegment, Iterable<ResourceAction>> SEGMENT_RA_GENERATOR = segment -> Collections.singletonList(AuthorizationUtils.DATASOURCE_READ_RA_GENERATOR.apply( segment.getDataSource()) ); - private static final long REPLICATION_FACTOR_UNKNOWN = -1L; - - /** - * Booleans constants represented as long type, - * where 1 = true and 0 = false to make it easy to count number of segments - * which are published, available etc. - */ - private static final long IS_ACTIVE_FALSE = 0L; - private static final long IS_ACTIVE_TRUE = 1L; - private static final long IS_PUBLISHED_FALSE = 0L; - private static final long IS_PUBLISHED_TRUE = 1L; - private static final long IS_AVAILABLE_TRUE = 1L; - private static final long IS_OVERSHADOWED_FALSE = 0L; - private static final long IS_OVERSHADOWED_TRUE = 1L; - public static boolean canUseNativeSystemTable( final RelOptTable table, final PlannerContext plannerContext @@ -157,47 +129,7 @@ public class SystemSchema extends AbstractTableSchema return nativeSystemTable == null ? null : nativeSystemTable.asNativeTable(); } - static final RowSignature SEGMENTS_SIGNATURE = RowSignature - .builder() - .add("segment_id", ColumnType.STRING) - .add("datasource", ColumnType.STRING) - .add("start", ColumnType.STRING) - .add("end", ColumnType.STRING) - .add("size", ColumnType.LONG) - .add("version", ColumnType.STRING) - .add("partition_num", ColumnType.LONG) - .add("num_replicas", ColumnType.LONG) - .add("num_rows", ColumnType.LONG) - .add("is_active", ColumnType.LONG) - .add("is_published", ColumnType.LONG) - .add("is_available", ColumnType.LONG) - .add("is_realtime", ColumnType.LONG) - .add("is_overshadowed", ColumnType.LONG) - .add("shard_spec", ColumnType.STRING) - .add("dimensions", ColumnType.STRING) - .add("metrics", ColumnType.STRING) - .add("projections", ColumnType.STRING) - .add("last_compaction_state", ColumnType.STRING) - .add("replication_factor", ColumnType.LONG) - .build(); - - /** - * List of [0..n) where n is the size of {@link #SEGMENTS_SIGNATURE}. - */ - private static final int[] SEGMENTS_PROJECT_ALL = IntStream.range(0, SEGMENTS_SIGNATURE.size()).toArray(); - - /** - * Fields in {@link #SEGMENTS_SIGNATURE} that are serialized with {@link ObjectMapper#writeValueAsString(Object)}. - */ - private static final IntSet SEGMENTS_JSON_FIELDS = new IntOpenHashSet( - new int[]{ - SEGMENTS_SIGNATURE.indexOf("shard_spec"), - SEGMENTS_SIGNATURE.indexOf("dimensions"), - SEGMENTS_SIGNATURE.indexOf("metrics"), - SEGMENTS_SIGNATURE.indexOf("projections"), - SEGMENTS_SIGNATURE.indexOf("last_compaction_state") - } - ); + static final RowSignature SEGMENTS_SIGNATURE = SegmentsTableDescriptor.ROW_SIGNATURE; static final RowSignature SERVERS_SIGNATURE = RowSignature .builder() @@ -407,7 +339,7 @@ public class SystemSchema extends AbstractTableSchema /** * This table contains row per segment from metadata store as well as served segments. */ - static class SegmentsTable extends AbstractTable implements ProjectableFilterableTable + static class SegmentsTable extends AbstractTable implements ProjectableFilterableTable, NativeSystemTable { private static final int DATASOURCE_COLUMN = SEGMENTS_SIGNATURE.indexOf("datasource"); @@ -444,6 +376,12 @@ public class SystemSchema extends AbstractTableSchema return TableType.SYSTEM_TABLE; } + @Override + public DruidTable asNativeTable() + { + return new NativeSegmentsTable(); + } + @Override public Enumerable<Object[]> scan( final DataContext root, @@ -451,164 +389,21 @@ public class SystemSchema extends AbstractTableSchema @Nullable final int[] projects ) { - // Best-effort push-down of a `datasource` equality/IN filter so we scan only the matching - // datasources instead of every segment in the cluster. Null => no usable filter => full scan. - // The filters are intentionally left in the list, so Calcite still applies them and correctness - // holds even if this extraction is conservative or over-broad. final Set<String> dataSourceFilter = getDataSourceFilter(filters); - - // Keep track of which segments we emitted from the publishedSegments iterator, so we don't emit them again - // from the availableSegments iterator. When a datasource filter is pushed down we only emit the matching - // datasources' segments, so avoid pre-sizing to the whole-cluster segment count (a huge, wasted allocation). - final Set<SegmentId> segmentsAlreadySeen = - dataSourceFilter == null - ? Sets.newHashSetWithExpectedSize(segmentMetadataCache.getTotalSegments()) - : new HashSet<>(); - - // Get segments from metadata segment cache (if enabled in SQL planner config), else directly from - // Coordinator. This may include both published and realtime segments. - final Iterator<SegmentStatusInCluster> metadataStoreSegments = metadataView.getSegments(dataSourceFilter); - final FluentIterable<Object[]> publishedSegments = FluentIterable - .from(() -> getAuthorizedPublishedSegments(metadataStoreSegments)) - .transform(val -> { - final DataSegment segment = val.getDataSegment(); - final AvailableSegmentMetadata availableSegmentMetadata = - segmentMetadataCache.getAvailableSegmentMetadata(segment.getDataSource(), segment.getId()); - segmentsAlreadySeen.add(segment.getId()); - - long numReplicas = 0L, isAvailable = 0L; - if (availableSegmentMetadata != null) { - numReplicas = availableSegmentMetadata.getNumReplicas(); - isAvailable = availableSegmentMetadata.getNumReplicas() > 0 ? IS_AVAILABLE_TRUE : IS_ACTIVE_FALSE; - } - - final long numRows; - if (segment.getTotalRows() != null) { - // the recent version of DataSegment stores numRows - numRows = segment.getTotalRows().longValue(); - } else if (val.getNumRows() != null) { - // If druid.centralizedDatasourceSchema.enabled is set on the Coordinator, SegmentMetadataCache on the - // broker might have outdated or no information regarding numRows and rowSignature for a segment. - // In that case, we should use {@code numRows} from the segment polled from the coordinator. - numRows = val.getNumRows(); - } else if (availableSegmentMetadata != null) { - numRows = availableSegmentMetadata.getNumRows(); - } else { - numRows = 0L; - } - - long isRealtime = val.isRealtime() ? 1 : 0; - - // set of segments returned from Coordinator include published and realtime segments - // so realtime segments are not published and vice versa - boolean isPublished = !val.isRealtime(); - - // is_active is true for published segments that are not overshadowed or else they should be realtime - boolean isActive = isPublished ? !val.isOvershadowed() : val.isRealtime(); - - return new Object[]{ - segment.getId(), - segment.getDataSource(), - segment.getInterval().getStart(), - segment.getInterval().getEnd(), - segment.getSize(), - segment.getVersion(), - (long) segment.getShardSpec().getPartitionNum(), - numReplicas, - numRows, - isActive ? IS_ACTIVE_TRUE : IS_ACTIVE_FALSE, - isPublished ? IS_PUBLISHED_TRUE : IS_PUBLISHED_FALSE, - isAvailable, - isRealtime, - val.isOvershadowed() ? IS_OVERSHADOWED_TRUE : IS_OVERSHADOWED_FALSE, - segment.getShardSpec(), - segment.getDimensions(), - segment.getMetrics(), - segment.getProjections(), - segment.getLastCompactionState(), - // If the segment is unpublished, we won't have this information yet. - // If the value is null, the load rules might have not evaluated yet, and we don't know the replication factor. - // This should be automatically updated in the next Coordinator poll. - val.getReplicationFactor() == null ? REPLICATION_FACTOR_UNKNOWN : (long) val.getReplicationFactor() - }; - }); - - // If druid.centralizedDatasourceSchema.enabled is set on the Coordinator, all the segments in this loop - // would be covered in the previous iteration since Coordinator would return realtime segments as well. - final FluentIterable<Object[]> availableSegments = FluentIterable - .from(() -> getAuthorizedAvailableSegments(segmentMetadataCache.iterateSegmentMetadata(dataSourceFilter))) - .transform(val -> { - final DataSegment segment = val.getSegment(); - if (segmentsAlreadySeen.contains(segment.getId())) { - return null; - } - return new Object[]{ - segment.getId(), - segment.getDataSource(), - segment.getInterval().getStart(), - segment.getInterval().getEnd(), - segment.getSize(), - segment.getVersion(), - (long) segment.getShardSpec().getPartitionNum(), - val.getNumReplicas(), - segment.getTotalRows() != null ? segment.getTotalRows() : val.getNumRows(), - // is_active is true for unpublished segments iff they are realtime - val.isRealtime() /* is_active */, - // is_published is false for unpublished segments - IS_PUBLISHED_FALSE, - // is_available is assumed to be always true for segments announced by historicals or realtime tasks - IS_AVAILABLE_TRUE, - val.isRealtime(), - IS_OVERSHADOWED_FALSE, - // there is an assumption here that unpublished segments are never overshadowed - segment.getShardSpec(), - segment.getDimensions(), - segment.getMetrics(), - segment.getProjections(), - null, // unpublished segments from realtime tasks will not be compacted yet - REPLICATION_FACTOR_UNKNOWN // If the segment is unpublished, we won't have this information yet. - }; - }); - - final Iterable<Object[]> allSegments = Iterables.unmodifiableIterable( - Iterables.concat(publishedSegments, availableSegments) + final Iterable<Object[]> authorizedRows = AuthorizationUtils.filterAuthorizedResources( + authenticationResult, + Iterables.filter( + SegmentsTableDataProvider.getRawRows(segmentMetadataCache, metadataView, dataSourceFilter), + Objects::nonNull + ), + row -> Collections.singletonList( + AuthorizationUtils.DATASOURCE_READ_RA_GENERATOR.apply((String) row[DATASOURCE_COLUMN]) + ), + authorizerMapper ); - - return Linq4j.asEnumerable(allSegments) + return Linq4j.asEnumerable(authorizedRows) .where(Objects::nonNull) - .select(row -> projectSegmentsRow(row, projects, jsonMapper)); - } - - private Iterator<SegmentStatusInCluster> getAuthorizedPublishedSegments(Iterator<SegmentStatusInCluster> it) - { - final Iterable<SegmentStatusInCluster> authorizedSegments = AuthorizationUtils - .filterAuthorizedResources( - authenticationResult, - () -> it, - SEGMENT_STATUS_IN_CLUSTER_RA_GENERATOR, - authorizerMapper - ); - return authorizedSegments.iterator(); - } - - private Iterator<AvailableSegmentMetadata> getAuthorizedAvailableSegments( - Iterator<AvailableSegmentMetadata> availableSegmentEntries - ) - { - Function<AvailableSegmentMetadata, Iterable<ResourceAction>> raGenerator = segment -> - Collections.singletonList( - AuthorizationUtils.DATASOURCE_READ_RA_GENERATOR.apply(segment.getSegment().getDataSource()) - ); - - final Iterable<AvailableSegmentMetadata> authorizedSegments = - AuthorizationUtils.filterAuthorizedResources( - authenticationResult, - () -> availableSegmentEntries, - raGenerator, - authorizerMapper - ); - - return authorizedSegments.iterator(); + .select(row -> SegmentsTableDataProvider.projectRow(row, projects, jsonMapper)); } /** @@ -625,47 +420,6 @@ public class SystemSchema extends AbstractTableSchema return SystemSchemaFilters.extractColumnValues(filters, DATASOURCE_COLUMN); } - private static class PartialSegmentData - { - private final long isAvailable; - private final long isRealtime; - private final long numReplicas; - private final long numRows; - - public PartialSegmentData( - final long isAvailable, - final long isRealtime, - final long numReplicas, - final long numRows - ) - - { - this.isAvailable = isAvailable; - this.isRealtime = isRealtime; - this.numReplicas = numReplicas; - this.numRows = numRows; - } - - public long isAvailable() - { - return isAvailable; - } - - public long isRealtime() - { - return isRealtime; - } - - public long getNumReplicas() - { - return numReplicas; - } - - public long getNumRows() - { - return numRows; - } - } } /** @@ -1286,46 +1040,6 @@ public class SystemSchema extends AbstractTableSchema .iterator(); } - /** - * Project a row using "projects" from {@link SegmentsTable#scan(DataContext, List, int[])}. - * <p> - * Also, fix up types so {@link ColumnType#STRING} are transformed to Strings if they aren't yet. This defers - * computation of {@link ObjectMapper#writeValueAsString(Object)} or {@link Object#toString()} until we know we - * actually need it. - */ - private static Object[] projectSegmentsRow( - final Object[] row, - @Nullable final int[] projects, - final ObjectMapper jsonMapper - ) - { - final int[] nonNullProjects = projects == null ? SEGMENTS_PROJECT_ALL : projects; - final Object[] projectedRow = new Object[nonNullProjects.length]; - - for (int i = 0; i < nonNullProjects.length; i++) { - final Object o = row[nonNullProjects[i]]; - - if (SEGMENTS_SIGNATURE.getColumnType(nonNullProjects[i]).get().is(ValueType.STRING) - && o != null - && !(o instanceof String)) { - // Delay calling toString() or ObjectMapper#writeValueAsString() until we know we actually need this field. - if (SEGMENTS_JSON_FIELDS.contains(nonNullProjects[i])) { - try { - projectedRow[i] = jsonMapper.writeValueAsString(o); - } - catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - } else { - projectedRow[i] = o.toString(); - } - } else { - projectedRow[i] = o; - } - } - return projectedRow; - } - /** * This table contains currently running and recently completed queries from all SQL engines. * Enabled based on {@link PlannerConfig#isEnableSysQueriesTable()}. 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 new file mode 100644 index 00000000000..8d8f54fa01f --- /dev/null +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SegmentsTableDataProviderTest.java @@ -0,0 +1,96 @@ +/* + * 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.sql.calcite.schema; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.inject.Provider; +import org.apache.druid.jackson.DefaultObjectMapper; +import org.apache.druid.java.util.common.Intervals; +import org.apache.druid.query.filter.DimFilter; +import org.apache.druid.query.filter.SelectorDimFilter; +import org.apache.druid.server.security.AuthConfig; +import org.apache.druid.server.security.AuthenticationResult; +import org.apache.druid.server.system.table.SegmentsTableDescriptor; +import org.apache.druid.timeline.DataSegment; +import org.apache.druid.timeline.SegmentId; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import java.util.Collections; +import java.util.List; +import java.util.Set; + +public class SegmentsTableDataProviderTest +{ + private static final AuthenticationResult AUTHENTICATION_RESULT = + new AuthenticationResult("test-user", AuthConfig.ALLOW_ALL_NAME, null, null); + + @Test + public void testDatasourceFilterIsPushedIntoBothLocalViews() + { + final BrokerSegmentMetadataCache metadataCache = Mockito.mock(BrokerSegmentMetadataCache.class); + final MetadataSegmentView metadataView = Mockito.mock(MetadataSegmentView.class); + Mockito.when(metadataView.getSegments(Set.of("foo"))).thenReturn(Collections.emptyIterator()); + Mockito.when(metadataCache.iterateSegmentMetadata(Set.of("foo"))).thenReturn(Collections.emptyIterator()); + + final SegmentsTableDataProvider provider = new SegmentsTableDataProvider( + (Provider<BrokerSegmentMetadataCache>) () -> metadataCache, + metadataView, + new DefaultObjectMapper() + ); + final List<DimFilter> filters = List.of(new SelectorDimFilter("datasource", "foo", null)); + + Assertions.assertTrue(toRows(provider.getRows(filters, AUTHENTICATION_RESULT)).isEmpty()); + Mockito.verify(metadataView).getSegments(Set.of("foo")); + Mockito.verify(metadataCache).iterateSegmentMetadata(Set.of("foo")); + Mockito.verify(metadataCache, Mockito.never()).getTotalSegments(); + } + + @Test + public void testComplexColumnsUseBindableJsonRepresentation() + { + final ObjectMapper mapper = new DefaultObjectMapper(); + final DataSegment segment = DataSegment.builder( + SegmentId.of("foo", Intervals.of("2000/2001"), "v", null) + ).dimensions(List.of("dim1")) + .metrics(List.of("metric1")) + .build(); + final Object[] rawRow = new Object[SegmentsTableDescriptor.ROW_SIGNATURE.size()]; + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("segment_id")] = segment.getId(); + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("datasource")] = segment.getDataSource(); + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("shard_spec")] = segment.getShardSpec(); + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("dimensions")] = segment.getDimensions(); + rawRow[SegmentsTableDescriptor.ROW_SIGNATURE.indexOf("metrics")] = segment.getMetrics(); + + final Object[] row = SegmentsTableDataProvider.projectRow(rawRow, null, mapper); + + Assertions.assertEquals(segment.getId().toString(), row[0]); + Assertions.assertEquals("[\"dim1\"]", row[15]); + Assertions.assertEquals("[\"metric1\"]", row[16]); + } + + private static List<Object[]> toRows(final Iterable<Object[]> rows) + { + final List<Object[]> result = new java.util.ArrayList<>(); + rows.forEach(result::add); + return result; + } +} diff --git a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemTableDataProviderTest.java b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemTableDataProviderTest.java index 4fb7561d84c..7ce3772ab8e 100644 --- a/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemTableDataProviderTest.java +++ b/sql/src/test/java/org/apache/druid/sql/calcite/schema/SystemTableDataProviderTest.java @@ -36,6 +36,12 @@ public class SystemTableDataProviderTest @Test public void testNativeTablesExposeSystemMetadata() { + final NativeSegmentsTable segments = new NativeSegmentsTable(); + Assertions.assertEquals("segments", ((SystemTableDataSource) segments.getDataSource()).getTable()); + Assertions.assertFalse(segments.isJoinable()); + Assertions.assertFalse(segments.isBroadcast()); + Assertions.assertEquals(Schema.TableType.SYSTEM_TABLE, segments.getJdbcTableType()); + final NativeServerPropertiesTable serverProperties = new NativeServerPropertiesTable(); Assertions.assertEquals( "server_properties", --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
