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 2a5a9de72843d12bfd63c2697fe6c22d8e963b01 Author: Frank Chen <[email protected]> AuthorDate: Fri Sep 11 14:29:09 2026 +0800 Add native sys.segments filter pushdown --- .../calcite/schema/SysSegmentsSqlBenchmark.java | 183 ++++++++++----------- docs/querying/sql-metadata-tables.md | 2 +- .../calcite/schema/SegmentsTableDataProvider.java | 61 ++++++- .../schema/SegmentsTableDataProviderTest.java | 96 +++++++++++ 4 files changed, 239 insertions(+), 103 deletions(-) diff --git a/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsSqlBenchmark.java b/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsSqlBenchmark.java index 05dc8dba786..21a99a23be8 100644 --- a/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsSqlBenchmark.java +++ b/benchmarks/src/test/java/org/apache/druid/sql/calcite/schema/SysSegmentsSqlBenchmark.java @@ -40,14 +40,11 @@ import org.apache.druid.query.BatchedInlineDataSource; import org.apache.druid.query.DataSource; import org.apache.druid.query.DefaultGenericQueryMetricsFactory; import org.apache.druid.query.DefaultQueryConfig; -import org.apache.druid.query.InlineDataSource; import org.apache.druid.query.QueryRunnerFactoryConglomerate; import org.apache.druid.query.SystemTableDataSource; -import org.apache.druid.query.filter.DimFilter; import org.apache.druid.query.policy.NoopPolicyEnforcer; import org.apache.druid.query.scan.ScanQueryEngine; import org.apache.druid.rpc.indexing.NoopOverlordClient; -import org.apache.druid.segment.InlineSegmentWrangler; import org.apache.druid.segment.MapSegmentWrangler; import org.apache.druid.segment.join.JoinableFactory; import org.apache.druid.segment.join.JoinableFactoryWrapper; @@ -61,7 +58,6 @@ import org.apache.druid.server.log.NoopRequestLogger; import org.apache.druid.server.metrics.NoopServiceEmitter; import org.apache.druid.server.security.AuthConfig; import org.apache.druid.server.security.AuthTestUtils; -import org.apache.druid.server.security.AuthenticationResult; import org.apache.druid.server.system.handler.SystemTableNodeLocator; import org.apache.druid.server.system.handler.SystemTableQueryClient; import org.apache.druid.server.system.handler.SystemTableQueryHandler; @@ -112,7 +108,11 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.TimeUnit; -/** Compares Bindable and native execution of the Web Console datasource-tab query over 500,000 segments. */ +/** + * Compares Bindable and batched native execution of the Web Console datasource-tab query over 500,000 segments. + * Filtered workloads also include a batched native control that disables provider pushdown while preserving the + * residual query filter, isolating the amount of work avoided by pushdown. + */ @State(Scope.Benchmark) @Fork( value = 1, @@ -134,35 +134,54 @@ public class SysSegmentsSqlBenchmark { private static final int NUM_SEGMENTS = 500_000; private static final int NUM_DATASOURCES = 1_000; + private static final String FIRST_SEGMENT_ID = SegmentId.of( + "datasource_0", + Intervals.utc(0, 86_400_000L), + "1", + new LinearShardSpec(0) + ).toString(); + + public enum Workload + { + SINGLE_DATASOURCE_EQUALS("WHERE datasource = 'datasource_0'\n", 1), + MULTI_DATASOURCE_IN("WHERE datasource IN ('datasource_0', 'datasource_1', 'datasource_2')\n", 3), + SEGMENT_ID_EQUALS("WHERE segment_id = '" + FIRST_SEGMENT_ID + "'\n", 1); + + private final String filterClause; + private final int expectedResultRows; + + Workload(final String filterClause, final int expectedResultRows) + { + this.filterClause = filterClause; + this.expectedResultRows = expectedResultRows; + } + } + private static final String BINDABLE = "bindable"; - private static final String NATIVE_ROW = "nativeRow"; private static final String NATIVE_PROVIDER = "nativeProvider"; - private static final String SQL = "SELECT\n" - + "datasource,\n" - + "COUNT(*) FILTER (WHERE is_active = 1) AS num_segments,\n" - + "COUNT(*) FILTER (WHERE is_published = 1 AND is_overshadowed = 0 " - + "AND replication_factor = 0) AS num_zero_replica_segments,\n" - + "COUNT(*) FILTER (WHERE is_published = 1 AND is_overshadowed = 0 " - + "AND is_available = 0 AND replication_factor > 0) AS num_segments_to_load,\n" - + "COUNT(*) FILTER (WHERE is_available = 1 AND is_active = 0) " - + "AS num_segments_to_drop,\n" - + "SUM(\"size\") FILTER (WHERE is_active = 1) AS total_data_size,\n" - + "MIN(\"num_rows\") FILTER (WHERE is_available = 1 AND is_realtime = 0) " - + "AS min_segment_rows,\n" - + "AVG(\"num_rows\") FILTER (WHERE is_available = 1 AND is_realtime = 0) " - + "AS avg_segment_rows,\n" - + "MAX(\"num_rows\") FILTER (WHERE is_available = 1 AND is_realtime = 0) " - + "AS max_segment_rows,\n" - + "SUM(\"num_rows\") FILTER (WHERE is_active = 1) AS total_rows,\n" - + "CASE WHEN SUM(\"num_rows\") FILTER (WHERE is_available = 1) <> 0 " - + "THEN (SUM(\"size\") FILTER (WHERE is_available = 1) / " - + "SUM(\"num_rows\") FILTER (WHERE is_available = 1)) ELSE 0 END " - + "AS avg_row_size,\n" - + "SUM(\"size\" * \"num_replicas\") FILTER (WHERE is_active = 1) " - + "AS replicated_size\n" - + "FROM sys.segments\n" - + "GROUP BY 1\n" - + "ORDER BY 1"; + private static final String NATIVE_PROVIDER_NO_PUSHDOWN = "nativeProviderNoPushdown"; + private static final String SQL_PREFIX = """ + SELECT + datasource, + COUNT(*) FILTER (WHERE is_active = 1) AS num_segments, + COUNT(*) FILTER (WHERE is_published = 1 AND is_overshadowed = 0 AND replication_factor = 0) + AS num_zero_replica_segments, + COUNT(*) FILTER (WHERE is_published = 1 AND is_overshadowed = 0 AND is_available = 0 AND replication_factor > 0) + AS num_segments_to_load, + COUNT(*) FILTER (WHERE is_available = 1 AND is_active = 0) AS num_segments_to_drop, + SUM("size") FILTER (WHERE is_active = 1) AS total_data_size, + MIN("num_rows") FILTER (WHERE is_available = 1 AND is_realtime = 0) AS min_segment_rows, + AVG("num_rows") FILTER (WHERE is_available = 1 AND is_realtime = 0) AS avg_segment_rows, + MAX("num_rows") FILTER (WHERE is_available = 1 AND is_realtime = 0) AS max_segment_rows, + SUM("num_rows") FILTER (WHERE is_active = 1) AS total_rows, + CASE WHEN SUM("num_rows") FILTER (WHERE is_available = 1) <> 0 + THEN (SUM("size") FILTER (WHERE is_available = 1) / SUM("num_rows") FILTER (WHERE is_available = 1)) + ELSE 0 + END AS avg_row_size, + SUM("size" * "num_replicas") FILTER (WHERE is_active = 1) AS replicated_size + FROM sys.segments + """; + private static final String SQL_SUFFIX = "GROUP BY 1\nORDER BY 1"; private static final Map<String, Object> BINDABLE_CONTEXT = ImmutableMap.of( PlannerContext.CTX_USE_NATIVE_QUERY_FOR_SYSTEM_TABLES, @@ -175,13 +194,16 @@ public class SysSegmentsSqlBenchmark private final Closer closer = Closer.create(); private PlannerFactory plannerFactory; - private SqlEngine rowEngine; private SqlEngine providerEngine; + private SqlEngine noPushdownProviderEngine; + + @Param + private Workload workload; @State(Scope.Thread) public static class ExecutionState { - @Param({BINDABLE, NATIVE_ROW, NATIVE_PROVIDER}) + @Param({BINDABLE, NATIVE_PROVIDER, NATIVE_PROVIDER_NO_PUSHDOWN}) private String executionPath; private PreparedQuery preparedQuery; @@ -237,47 +259,6 @@ public class SysSegmentsSqlBenchmark } } - /** Keeps the former row-only native path available as a stable benchmark baseline. */ - private static class RowOnlySystemTableDataProvider implements SystemTableDataProvider - { - private final SystemTableDataProvider delegate; - - RowOnlySystemTableDataProvider(final SystemTableDataProvider delegate) - { - this.delegate = delegate; - } - - @Override - public List<SystemTablePushdownFilter> getPushdownFilters() - { - return delegate.getPushdownFilters(); - } - - @Override - public Iterable<Object[]> getRows( - final List<DimFilter> filters, - final AuthenticationResult internalAuthenticationResult - ) - { - return delegate.getRows(filters, internalAuthenticationResult); - } - - @Override - public Iterable<Object[]> getRawRows( - final List<DimFilter> filters, - final AuthenticationResult internalAuthenticationResult - ) - { - return delegate.getRawRows(filters, internalAuthenticationResult); - } - - @Override - public Object[] projectRow(final Object[] row, final int[] projects) - { - return delegate.projectRow(row, projects); - } - } - @Setup(Level.Trial) public void setup() { @@ -310,6 +291,19 @@ public class SysSegmentsSqlBenchmark metadataView, CalciteTests.getJsonMapper() ); + final SegmentsTableDataProvider noPushdownDataProvider = new SegmentsTableDataProvider( + () -> segmentMetadataCache, + metadataView, + CalciteTests.getJsonMapper() + ) + { + /** Keeps all other provider behavior identical while preventing provider-level filtering. */ + @Override + public List<SystemTablePushdownFilter> getPushdownFilters() + { + return Collections.emptyList(); + } + }; final QueryRunnerFactoryConglomerate conglomerate = QueryStackTests.createQueryRunnerFactoryConglomerate(closer); final SpecificSegmentsQuerySegmentWalker walker = closer.register( @@ -318,8 +312,6 @@ public class SysSegmentsSqlBenchmark conglomerate, new MapSegmentWrangler( Map.of( - InlineDataSource.class, - new InlineSegmentWrangler(), BatchedInlineDataSource.class, new BatchedInlineDataSource.Wrangler() ) @@ -334,17 +326,14 @@ public class SysSegmentsSqlBenchmark new ScanQueryEngine(), AuthTestUtils.TEST_AUTHORIZER_MAPPER ); - final SystemTableQueryHandler rowQueryHandler = new SystemTableQueryHandler( - Map.<String, SystemTableDataProvider>of( - descriptor.getTableName(), - new RowOnlySystemTableDataProvider(dataProvider) - ), + final SystemTableQueryHandler noPushdownProviderQueryHandler = new SystemTableQueryHandler( + Map.<String, SystemTableDataProvider>of(descriptor.getTableName(), noPushdownDataProvider), Map.<String, SystemTableDescriptor>of(descriptor.getTableName(), descriptor), new ScanQueryEngine(), AuthTestUtils.TEST_AUTHORIZER_MAPPER ); - rowEngine = makeEngine(conglomerate, walker, descriptor, rowQueryHandler); providerEngine = makeEngine(conglomerate, walker, descriptor, providerQueryHandler); + noPushdownProviderEngine = makeEngine(conglomerate, walker, descriptor, noPushdownProviderQueryHandler); final PlannerConfig plannerConfig = new PlannerConfig(); final TimelineServerView timelineServerView = new TestTimelineServerView(Collections.emptyList()); @@ -399,9 +388,9 @@ public class SysSegmentsSqlBenchmark ); final List<Object[]> bindableResults = runQuery(BINDABLE); - for (final String executionPath : List.of(NATIVE_ROW, NATIVE_PROVIDER)) { + for (final String executionPath : List.of(NATIVE_PROVIDER, NATIVE_PROVIDER_NO_PUSHDOWN)) { final List<Object[]> nativeResults = runQuery(executionPath); - if (bindableResults.size() != NUM_DATASOURCES || !rowsEqual(bindableResults, nativeResults)) { + if (bindableResults.size() != workload.expectedResultRows || !rowsEqual(bindableResults, nativeResults)) { throw new IllegalStateException("Bindable and native benchmark results do not match for " + executionPath); } } @@ -482,22 +471,26 @@ public class SysSegmentsSqlBenchmark final Map<String, Object> context; switch (executionPath) { case BINDABLE: - engine = rowEngine; + engine = providerEngine; context = BINDABLE_CONTEXT; break; - case NATIVE_ROW: - engine = rowEngine; - context = NATIVE_CONTEXT; - break; case NATIVE_PROVIDER: engine = providerEngine; context = NATIVE_CONTEXT; break; + case NATIVE_PROVIDER_NO_PUSHDOWN: + engine = noPushdownProviderEngine; + context = NATIVE_CONTEXT; + break; default: throw new IllegalArgumentException("Unknown execution path " + executionPath); } - final DruidPlanner planner = plannerFactory.createPlannerForTesting(engine, SQL, context); + final DruidPlanner planner = plannerFactory.createPlannerForTesting( + engine, + SQL_PREFIX + workload.filterClause + SQL_SUFFIX, + context + ); try { return new PreparedQuery(planner, planner.plan()); } @@ -521,15 +514,15 @@ public class SysSegmentsSqlBenchmark } @Benchmark - public void queryNative(final Blackhole blackhole) + public void queryNativeProvider(final Blackhole blackhole) { - blackhole.consume(runQuery(NATIVE_ROW)); + blackhole.consume(runQuery(NATIVE_PROVIDER)); } @Benchmark - public void queryNativeProvider(final Blackhole blackhole) + public void queryNativeProviderWithoutPushdown(final Blackhole blackhole) { - blackhole.consume(runQuery(NATIVE_PROVIDER)); + blackhole.consume(runQuery(NATIVE_PROVIDER_NO_PUSHDOWN)); } @Benchmark diff --git a/docs/querying/sql-metadata-tables.md b/docs/querying/sql-metadata-tables.md index 61c40d5dede..fd51c3e5517 100644 --- a/docs/querying/sql-metadata-tables.md +++ b/docs/querying/sql-metadata-tables.md @@ -168,7 +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.segments`](#segments-table)|The Broker that receives the SQL query. It combines that Broker's segment metadata cache with its Coordinator metadata view. Exact equality/`IN` filters on `datasource`, `segment_id`, `start`, `end`, and `version` are applied before matching rows are built; `datasource` also restricts both local metadata views. 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/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 923be45bf4a..267c42414cb 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 @@ -48,12 +48,15 @@ import org.apache.druid.timeline.SegmentId; import org.apache.druid.timeline.SegmentStatusInCluster; import javax.annotation.Nullable; +import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.Set; +import java.util.function.Predicate; import java.util.stream.IntStream; /** Native row supplier for {@code sys.segments}. */ @@ -79,7 +82,11 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider } ); private static final List<SystemTablePushdownFilter> PUSHDOWN_FILTERS = List.of( - new SystemTablePushdownFilter("datasource", null) + new SystemTablePushdownFilter("datasource", null), + new SystemTablePushdownFilter("segment_id", null), + new SystemTablePushdownFilter("start", null), + new SystemTablePushdownFilter("end", null), + new SystemTablePushdownFilter("version", null) ); private final Provider<BrokerSegmentMetadataCache> segmentMetadataCacheProvider; @@ -139,8 +146,14 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider final AuthenticationResult internalAuthenticationResult ) { + final Map<String, Set<String>> stringFilters = getStringFilters(filters); return Iterables.filter( - getRawRows(segmentMetadataCacheProvider.get(), metadataView, getDataSourceFilter(filters)), + getRawRows( + segmentMetadataCacheProvider.get(), + metadataView, + stringFilters.get("datasource"), + segment -> matchesStringFilters(segment, stringFilters) + ), Objects::nonNull ); } @@ -157,6 +170,16 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider final MetadataSegmentView metadataView, @Nullable final Set<String> dataSourceFilter ) + { + return getRawRows(segmentMetadataCache, metadataView, dataSourceFilter, segment -> true); + } + + private static Iterable<Object[]> getRawRows( + final BrokerSegmentMetadataCache segmentMetadataCache, + final MetadataSegmentView metadataView, + @Nullable final Set<String> dataSourceFilter, + final Predicate<DataSegment> segmentFilter + ) { final Set<SegmentId> segmentsAlreadySeen = dataSourceFilter == null ? Sets.newHashSetWithExpectedSize( @@ -167,6 +190,7 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider final Iterator<SegmentStatusInCluster> metadataStoreSegments = metadataView.getSegments(dataSourceFilter); final FluentIterable<Object[]> publishedSegments = FluentIterable .from(() -> metadataStoreSegments) + .filter(val -> segmentFilter.test(val.getDataSegment())) .transform(val -> { final DataSegment segment = val.getDataSegment(); final AvailableSegmentMetadata availableSegmentMetadata = @@ -221,6 +245,7 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider final FluentIterable<Object[]> availableSegments = FluentIterable .from(() -> segmentMetadataCache.iterateSegmentMetadata(dataSourceFilter)) + .filter(val -> segmentFilter.test(val.getSegment())) .transform(val -> { final DataSegment segment = val.getSegment(); if (segmentsAlreadySeen.contains(segment.getId())) { @@ -286,19 +311,41 @@ public class SegmentsTableDataProvider implements SystemTableDataProvider return projectedRow; } - @Nullable - private static Set<String> getDataSourceFilter(final List<DimFilter> filters) + private static Map<String, Set<String>> getStringFilters(final List<DimFilter> filters) { - Set<String> result = null; + final Map<String, Set<String>> result = new HashMap<>(); for (final DimFilter filter : filters) { - if (isStringValuesFilter(filter) && "datasource".equals(SystemTablePushdownFilter.getStringValuesColumn(filter))) { + if (isStringValuesFilter(filter)) { + final String column = SystemTablePushdownFilter.getStringValuesColumn(filter); final Set<String> values = SystemTablePushdownFilter.getStringValues(filter); - result = result == null ? values : Sets.intersection(result, values).immutableCopy(); + result.compute( + column, + (ignored, currentValues) -> currentValues == null + ? values + : Sets.intersection(currentValues, values).immutableCopy() + ); } } return result; } + private static boolean matchesStringFilters( + final DataSegment segment, + final Map<String, Set<String>> filters + ) + { + return matches(filters.get("datasource"), segment.getDataSource()) + && matches(filters.get("segment_id"), segment.getId().toString()) + && matches(filters.get("start"), segment.getInterval().getStart().toString()) + && matches(filters.get("end"), segment.getInterval().getEnd().toString()) + && matches(filters.get("version"), segment.getVersion()); + } + + private static boolean matches(@Nullable final Set<String> acceptedValues, final String value) + { + return acceptedValues == null || acceptedValues.contains(value); + } + private static boolean isStringValuesFilter(final DimFilter filter) { if (filter instanceof SelectorDimFilter 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 564444cb693..4aee1d3a71d 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 @@ -29,18 +29,21 @@ import org.apache.druid.query.filter.DimFilter; import org.apache.druid.query.filter.SelectorDimFilter; import org.apache.druid.segment.column.ColumnType; import org.apache.druid.segment.column.RowSignature; +import org.apache.druid.segment.metadata.AvailableSegmentMetadata; 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.server.system.table.SystemTableQueryRequest; import org.apache.druid.timeline.DataSegment; import org.apache.druid.timeline.SegmentId; +import org.apache.druid.timeline.SegmentStatusInCluster; 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.Map; import java.util.Set; public class SegmentsTableDataProviderTest @@ -48,6 +51,21 @@ public class SegmentsTableDataProviderTest private static final AuthenticationResult AUTHENTICATION_RESULT = new AuthenticationResult("test-user", AuthConfig.ALLOW_ALL_NAME, null, null); + @Test + public void testPushdownFilters() + { + final SegmentsTableDataProvider provider = new SegmentsTableDataProvider( + () -> Mockito.mock(BrokerSegmentMetadataCache.class), + Mockito.mock(MetadataSegmentView.class), + new DefaultObjectMapper() + ); + + Assertions.assertEquals( + List.of("datasource", "segment_id", "start", "end", "version"), + provider.getPushdownFilters().stream().map(filter -> filter.key()).toList() + ); + } + @Test public void testDatasourceFilterIsPushedIntoBothLocalViews() { @@ -69,6 +87,84 @@ public class SegmentsTableDataProviderTest Mockito.verify(metadataCache, Mockito.never()).getTotalSegments(); } + @Test + public void testExactSegmentFieldsAreFilteredBeforeRowsAreBuilt() + { + final DataSegment matchingSegment = DataSegment.builder( + SegmentId.of("foo", Intervals.of("2000/2001"), "v1", null) + ).build(); + final DataSegment otherSegment = DataSegment.builder( + SegmentId.of("foo", Intervals.of("2001/2002"), "v2", null) + ).build(); + final List<SegmentStatusInCluster> segments = List.of( + new SegmentStatusInCluster(matchingSegment, false, 1, 10L, false), + new SegmentStatusInCluster(otherSegment, false, 1, 10L, false) + ); + final BrokerSegmentMetadataCache metadataCache = Mockito.mock(BrokerSegmentMetadataCache.class); + final MetadataSegmentView metadataView = Mockito.mock(MetadataSegmentView.class); + Mockito.when(metadataView.getSegments((Set<String>) null)).thenAnswer(ignored -> segments.iterator()); + Mockito.when(metadataCache.iterateSegmentMetadata(null)).thenAnswer(ignored -> Collections.emptyIterator()); + Mockito.when(metadataCache.getTotalSegments()).thenReturn(segments.size()); + + final SegmentsTableDataProvider provider = new SegmentsTableDataProvider( + () -> metadataCache, + metadataView, + new DefaultObjectMapper() + ); + final Map<String, String> filters = Map.of( + "segment_id", matchingSegment.getId().toString(), + "start", matchingSegment.getInterval().getStart().toString(), + "end", matchingSegment.getInterval().getEnd().toString(), + "version", matchingSegment.getVersion() + ); + + for (final Map.Entry<String, String> filter : filters.entrySet()) { + final List<Object[]> rows = toRows( + provider.getRows( + List.of(new SelectorDimFilter(filter.getKey(), filter.getValue(), null)), + AUTHENTICATION_RESULT + ) + ); + Assertions.assertEquals(1, rows.size(), filter.getKey()); + Assertions.assertEquals(matchingSegment.getId().toString(), rows.get(0)[0], filter.getKey()); + } + } + + @Test + public void testExactSegmentFieldsFilterAvailableSegments() + { + final DataSegment matchingSegment = DataSegment.builder( + SegmentId.of("foo", Intervals.of("2000/2001"), "v1", null) + ).build(); + final DataSegment otherSegment = DataSegment.builder( + SegmentId.of("foo", Intervals.of("2001/2002"), "v2", null) + ).build(); + final List<AvailableSegmentMetadata> segments = List.of( + AvailableSegmentMetadata.builder(matchingSegment, 0, Collections.emptySet(), null, 10).build(), + AvailableSegmentMetadata.builder(otherSegment, 0, Collections.emptySet(), null, 10).build() + ); + final BrokerSegmentMetadataCache metadataCache = Mockito.mock(BrokerSegmentMetadataCache.class); + final MetadataSegmentView metadataView = Mockito.mock(MetadataSegmentView.class); + Mockito.when(metadataView.getSegments((Set<String>) null)).thenAnswer(ignored -> Collections.emptyIterator()); + Mockito.when(metadataCache.iterateSegmentMetadata(null)).thenAnswer(ignored -> segments.iterator()); + Mockito.when(metadataCache.getTotalSegments()).thenReturn(segments.size()); + + final SegmentsTableDataProvider provider = new SegmentsTableDataProvider( + () -> metadataCache, + metadataView, + new DefaultObjectMapper() + ); + final List<Object[]> rows = toRows( + provider.getRows( + List.of(new SelectorDimFilter("segment_id", matchingSegment.getId().toString(), null)), + AUTHENTICATION_RESULT + ) + ); + + Assertions.assertEquals(1, rows.size()); + Assertions.assertEquals(matchingSegment.getId().toString(), rows.get(0)[0]); + } + @Test public void testComplexColumnsUseBindableJsonRepresentation() { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
