This is an automated email from the ASF dual-hosted git repository. xiangfu0 pushed a commit to branch xiangfu0/data-3221-12-parallel-value-pruning in repository https://gitbox.apache.org/repos/asf/pinot.git
commit 8c6e53471d5b325030c3e24431bb7ae1bdf49128 Author: Xiang Fu <[email protected]> AuthorDate: Tue Sep 8 17:14:49 2026 -0700 DATA-3221 (14): prune values across the query executor `ValueBasedSegmentPruner#prune` is a serial loop over every segment the server holds, and it runs before any per-segment parallelism starts, so on a server holding tens of thousands of segments it is a query's longest single-threaded stretch. It implemented only the two-arg `prune`, so `SegmentPruner`'s three-arg default silently discarded the `ExecutorService` it was handed. `BloomFilterSegmentPruner` already overrode that; the pruner that runs first, over the full segment set, did not. The parallel path mirrors the bloom-filter one: tasks over a strided slice of the segments, each with its own value and data-source caches — neither is thread-safe, and both are scoped to a single segment at a time regardless — run through `QueryMultiThreadingUtils.runTasksWithDeadline` so the query deadline and cancellation are honored. Below `TARGET_NUM_SEGMENTS_PER_THREAD` segments, or with no executor, it stays serial. Segments come back in a different order than they went in, which is already true of the bloom-filter pruner that runs immediately after this one. `TARGET_NUM_SEGMENTS_PER_THREAD` moves to the parent, where both pruners now read the same value instead of the child shadowing it. Co-Authored-By: Claude Opus 5 <[email protected]> --- .../query/pruner/BloomFilterSegmentPruner.java | 1 - .../core/query/pruner/ValueBasedSegmentPruner.java | 53 ++++++++++++++++++++++ .../query/pruner/ColumnValueSegmentPrunerTest.java | 47 +++++++++++++++++++ 3 files changed, 100 insertions(+), 1 deletion(-) diff --git a/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/BloomFilterSegmentPruner.java b/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/BloomFilterSegmentPruner.java index 44b7d99a784..6de1c9098b2 100644 --- a/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/BloomFilterSegmentPruner.java +++ b/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/BloomFilterSegmentPruner.java @@ -51,7 +51,6 @@ public class BloomFilterSegmentPruner extends ValueBasedSegmentPruner { // Try to schedule 10 segments for each thread, or evenly distribute them to all MAX_NUM_THREADS_PER_QUERY threads. // TODO: make this threshold configurable? threshold 10 is also used in CombinePlanNode, which accesses the // dictionary data to do query planning and if segments are more than 10, planning is done in parallel. - private static final int TARGET_NUM_SEGMENTS_PER_THREAD = 10; private FetchPlanner _fetchPlanner; diff --git a/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/ValueBasedSegmentPruner.java b/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/ValueBasedSegmentPruner.java index f5f73d2cfcc..091b42c7e75 100644 --- a/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/ValueBasedSegmentPruner.java +++ b/pinot-core/src/main/java/org/apache/pinot/core/query/pruner/ValueBasedSegmentPruner.java @@ -24,6 +24,8 @@ import java.util.IdentityHashMap; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.concurrent.ExecutorService; +import javax.annotation.Nullable; import org.apache.pinot.common.request.context.ExpressionContext; import org.apache.pinot.common.request.context.FilterContext; import org.apache.pinot.common.request.context.predicate.EqPredicate; @@ -31,6 +33,7 @@ import org.apache.pinot.common.request.context.predicate.InPredicate; import org.apache.pinot.common.request.context.predicate.Predicate; import org.apache.pinot.common.utils.config.QueryOptionsUtils; import org.apache.pinot.core.query.request.context.QueryContext; +import org.apache.pinot.core.util.QueryMultiThreadingUtils; import org.apache.pinot.segment.local.segment.index.readers.bloom.GuavaBloomFilterReaderUtils; import org.apache.pinot.segment.spi.IndexSegment; import org.apache.pinot.segment.spi.datasource.DataSource; @@ -38,6 +41,7 @@ import org.apache.pinot.segment.spi.index.reader.BloomFilterReader; import org.apache.pinot.spi.data.FieldSpec.DataType; import org.apache.pinot.spi.env.PinotConfiguration; import org.apache.pinot.spi.exception.BadQueryRequestException; +import org.apache.pinot.spi.exception.QueryCancelledException; import org.apache.pinot.spi.utils.CommonConstants.Server; @@ -108,6 +112,55 @@ abstract public class ValueBasedSegmentPruner implements SegmentPruner { abstract boolean isApplicableToPredicate(Predicate predicate, Map<String, String> queryOptions); + /// Prunes across the query executor when there are enough segments to be worth it. + /// + /// The serial [#prune(List, QueryContext)] below runs on the calling thread for every segment the server holds, + /// before any per-segment parallelism starts, so on a server holding tens of thousands of segments it is the + /// query's longest single-threaded stretch. [BloomFilterSegmentPruner] already overrides this; the pruner that + /// runs first, over the full segment set, did not, and [SegmentPruner#prune(List, QueryContext, ExecutorService)] + /// silently discarded the executor for it. + /// + /// Each task keeps its own value and data-source caches, as the parallel bloom-filter path does — neither is + /// thread-safe, and both are scoped to one segment at a time anyway. Segments come back in a different order than + /// they went in, which is already true of the bloom-filter pruner that runs immediately after this one. + protected static final int TARGET_NUM_SEGMENTS_PER_THREAD = 10; + + @Override + public List<IndexSegment> prune(List<IndexSegment> segments, QueryContext query, + @Nullable ExecutorService executorService) { + if (executorService == null || segments.size() <= TARGET_NUM_SEGMENTS_PER_THREAD) { + return prune(segments, query); + } + int numSegments = segments.size(); + int numTasks = QueryMultiThreadingUtils.getNumTasks(numSegments, TARGET_NUM_SEGMENTS_PER_THREAD, + query.getMaxExecutionThreads()); + List<IndexSegment> allSelectedSegments = new ArrayList<>(numSegments); + QueryMultiThreadingUtils.runTasksWithDeadline(numTasks, index -> { + FilterContext filter = Objects.requireNonNull(query.getFilter()); + ValueCache cachedValues = new ValueCache(); + Map<String, DataSource> dataSourceCache = new HashMap<>(); + List<IndexSegment> selectedSegments = new ArrayList<>(); + for (int i = index; i < numSegments; i += numTasks) { + dataSourceCache.clear(); + IndexSegment segment = segments.get(i); + if (!pruneSegment(segment, filter, dataSourceCache, cachedValues, query)) { + selectedSegments.add(segment); + } + } + return selectedSegments; + }, taskRes -> { + if (taskRes != null) { + allSelectedSegments.addAll(taskRes); + } + }, e -> { + if (e instanceof InterruptedException) { + throw new QueryCancelledException("Cancelled while running " + getClass().getSimpleName(), e); + } + throw new RuntimeException("Caught exception while running " + getClass().getSimpleName(), e); + }, executorService, query.getEndTimeMs()); + return allSelectedSegments; + } + @Override public List<IndexSegment> prune(List<IndexSegment> segments, QueryContext query) { if (segments.isEmpty()) { diff --git a/pinot-core/src/test/java/org/apache/pinot/core/query/pruner/ColumnValueSegmentPrunerTest.java b/pinot-core/src/test/java/org/apache/pinot/core/query/pruner/ColumnValueSegmentPrunerTest.java index c5cd96ec014..e8b21bd48d3 100644 --- a/pinot-core/src/test/java/org/apache/pinot/core/query/pruner/ColumnValueSegmentPrunerTest.java +++ b/pinot-core/src/test/java/org/apache/pinot/core/query/pruner/ColumnValueSegmentPrunerTest.java @@ -19,10 +19,15 @@ package org.apache.pinot.core.query.pruner; import com.google.common.collect.ImmutableSet; +import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; +import java.util.HashSet; +import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import org.apache.pinot.core.query.request.context.QueryContext; import org.apache.pinot.core.query.request.context.utils.QueryContextConverterUtils; import org.apache.pinot.segment.spi.IndexSegment; @@ -40,6 +45,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; @@ -233,6 +239,47 @@ public class ColumnValueSegmentPrunerTest { assertTrue(PRUNER.isApplicableTo(queryContext)); } + /// The pruner runs over every segment the server holds, so above a threshold it prunes across the query executor. + /// The parallel path must select exactly the same segments as the serial one. + @Test + public void testParallelPruningSelectsTheSameSegments() throws Exception { + int numSegments = 40; + List<IndexSegment> segments = new ArrayList<>(numSegments); + for (int i = 0; i < numSegments; i++) { + // Alternate: half the segments hold values the predicate can match, half cannot and must be pruned. + segments.add(segmentWithRange(i % 2 == 0 ? 0 : 100, i % 2 == 0 ? 50 : 150)); + } + QueryContext serialQuery = QueryContextConverterUtils.getQueryContext( + "SELECT COUNT(*) FROM testTable WHERE column = 10"); + serialQuery.setSchema(mock(Schema.class)); + QueryContext parallelQuery = QueryContextConverterUtils.getQueryContext( + "SELECT COUNT(*) FROM testTable WHERE column = 10"); + parallelQuery.setSchema(mock(Schema.class)); + parallelQuery.setEndTimeMs(System.currentTimeMillis() + 30_000); + + List<IndexSegment> serial = PRUNER.prune(segments, serialQuery); + ExecutorService executor = Executors.newFixedThreadPool(4); + try { + List<IndexSegment> parallel = PRUNER.prune(segments, parallelQuery, executor); + assertEquals(new HashSet<>(parallel), new HashSet<>(serial)); + assertEquals(parallel.size(), numSegments / 2); + } finally { + executor.shutdownNow(); + } + } + + private IndexSegment segmentWithRange(int minValue, int maxValue) { + IndexSegment indexSegment = mockIndexSegment(); + DataSource dataSource = mock(DataSource.class); + when(indexSegment.getDataSource(eq("column"), any(Schema.class))).thenReturn(dataSource); + DataSourceMetadata metadata = mock(DataSourceMetadata.class); + when(metadata.getDataType()).thenReturn(DataType.INT); + when(metadata.getMinValue()).thenReturn(minValue); + when(metadata.getMaxValue()).thenReturn(maxValue); + when(dataSource.getDataSourceMetadata()).thenReturn(metadata); + return indexSegment; + } + private IndexSegment mockIndexSegment() { IndexSegment indexSegment = mock(IndexSegment.class); when(indexSegment.getColumnNames()).thenReturn(ImmutableSet.of("column")); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
