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]

Reply via email to