Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
38 commits
Select commit Hold shift + click to select a range
469cf51
Combining filter rewrite and skip list approaches for further optimiz…
jainankitk Oct 8, 2025
e20f702
Removing parent aggregation check for perf benchmark
jainankitk Oct 8, 2025
82bc95d
Adding changelog entry
jainankitk Oct 8, 2025
aff3dc6
Applying the skip list optimization for AutoDateHistogram
jainankitk Oct 9, 2025
1c29540
Addressing checkstyle failures
jainankitk Oct 9, 2025
b9e9f2b
Apply spotless
jainankitk Oct 9, 2025
a28b9c1
Merge branch 'main' into agg-perf
jainankitk Oct 9, 2025
0a9ef40
Minor bug fix
jainankitk Oct 10, 2025
8d4ccf7
Merge branch 'main' into agg-perf
jainankitk Oct 21, 2025
3cd5f64
Updating to lucene 10.4 snapshot
jainankitk Oct 20, 2025
2c657e7
Using bulk collect APIs in Lucene
jainankitk Oct 21, 2025
879bfaf
Apply spotless
jainankitk Oct 21, 2025
128b9b4
Fixing build issues after changing to 10.4 lucene snapshot
jainankitk Oct 21, 2025
65192fd
Fixing Lucene103Codec references in test files
jainankitk Oct 21, 2025
931614c
Updating Lucene version for ES
jainankitk Oct 21, 2025
866937d
Merge branch 'main' into agg-perf
jainankitk Oct 21, 2025
2d697fc
Updating Lucene codec version for ES
jainankitk Oct 21, 2025
e04a1b4
Adding CompositeCodec104 changes
jainankitk Oct 22, 2025
45512cd
Remaining CompositeCodec104 changes
jainankitk Oct 22, 2025
729bf0a
Final CompositeCodec104 changes
jainankitk Oct 22, 2025
4e7d9e6
Merge branch 'main' into agg-perf
jainankitk Oct 23, 2025
09d5c71
Merge branch 'feature/3.x-lucene' into agg-perf-lucene
jainankitk Oct 23, 2025
23fbad3
Add unit test for filter rewrite with date histogram with skiplist.
asimmahmood1 Oct 27, 2025
7eb64f7
Spotless check
asimmahmood1 Oct 27, 2025
35834e4
Fix unit test
asimmahmood1 Oct 27, 2025
2b593c9
Merge remote-tracking branch 'upstream/main' into agg-perf
asimmahmood1 Oct 27, 2025
3cdc37d
Not ready for check-in, just throwing this out to come up with differ…
asimmahmood1 Nov 10, 2025
d0eeb37
Revert auto date changes for this PR
asimmahmood1 Nov 10, 2025
66ffef1
Merge remote-tracking branch 'upstream/main' into agg-perf
asimmahmood1 Nov 10, 2025
0ec357a
Switch to Lucene's version of BitSetDocIdStream
asimmahmood1 Nov 12, 2025
7a7209f
Merge remote-tracking branch 'upstream/main' into agg-perf
asimmahmood1 Nov 12, 2025
37f4641
Merge branch 'main' into agg-perf
jainankitk Nov 14, 2025
eaf7e52
Resolving merge conflict issue
jainankitk Nov 14, 2025
4d9dd28
Merge branch 'feature/3.x-lucene' into agg-perf-lucene
jainankitk Nov 14, 2025
8403592
Merge branch 'agg-perf' into agg-perf-lucene
jainankitk Nov 14, 2025
80079be
Fixing build failure
jainankitk Nov 14, 2025
e90eaca
Merge branch 'feature/3.x-lucene' into agg-perf-lucene
jainankitk Nov 14, 2025
17c6b7a
Merge branch 'feature/3.x-lucene' into agg-perf-lucene
jainankitk Nov 15, 2025
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
- Support pull-based ingestion message mappers and raw payload support ([#19765](https://github.com/opensearch-project/OpenSearch/pull/19765)]

### Changed
- Combining filter rewrite and skip list to optimize sub aggregation([#19573](https://github.com/opensearch-project/OpenSearch/pull/19573))
- Faster `terms` query creation for `keyword` field with index and docValues enabled ([#19350](https://github.com/opensearch-project/OpenSearch/pull/19350))
- Refactor to move prepareIndex and prepareDelete methods to Engine class ([#19551](https://github.com/opensearch-project/OpenSearch/pull/19551))
- Omit maxScoreCollector in SimpleTopDocsCollectorContext when concurrent segment search enabled ([#19584](https://github.com/opensearch-project/OpenSearch/pull/19584))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,12 @@ public void collect(int doc) throws IOException {
collect(doc, 0);
}

public void collect(int[] docIds, long owningBucketOrd) throws IOException {
for (int doc : docIds) {
collect(doc, owningBucketOrd);
}
}

@Override
public void collect(DocIdStream stream) throws IOException {
collect(stream, 0);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
/*
* SPDX-License-Identifier: Apache-2.0
*
* The OpenSearch Contributors require contributions made to
* this file be licensed under the Apache-2.0 license or a
* compatible open source license.
*/

package org.opensearch.search.aggregations.bucket;

import org.apache.lucene.index.DocValuesSkipper;
import org.apache.lucene.index.NumericDocValues;
import org.apache.lucene.search.DocIdStream;
import org.apache.lucene.search.Scorable;
import org.opensearch.common.Rounding;
import org.opensearch.search.aggregations.LeafBucketCollector;
import org.opensearch.search.aggregations.bucket.terms.LongKeyedBucketOrds;

import java.io.IOException;

/**
* Histogram collection logic using skip list.
*
* @opensearch.internal
*/
public class HistogramSkiplistLeafCollector extends LeafBucketCollector {

private final NumericDocValues values;
private final DocValuesSkipper skipper;
private final Rounding.Prepared preparedRounding;
private final LongKeyedBucketOrds bucketOrds;
private final LeafBucketCollector sub;
private final BucketsAggregator aggregator;

/**
* Max doc ID (inclusive) up to which all docs values may map to the same
* bucket.
*/
private int upToInclusive = -1;

/**
* Whether all docs up to {@link #upToInclusive} values map to the same bucket.
*/
private boolean upToSameBucket;

/**
* Index in bucketOrds for docs up to {@link #upToInclusive}.
*/
private long upToBucketIndex;

public HistogramSkiplistLeafCollector(
NumericDocValues values,
DocValuesSkipper skipper,
Rounding.Prepared preparedRounding,
LongKeyedBucketOrds bucketOrds,
LeafBucketCollector sub,
BucketsAggregator aggregator
) {
this.values = values;
this.skipper = skipper;
this.preparedRounding = preparedRounding;
this.bucketOrds = bucketOrds;
this.sub = sub;
this.aggregator = aggregator;
}

@Override
public void setScorer(Scorable scorer) throws IOException {
if (sub != null) {
sub.setScorer(scorer);
}
}

private void advanceSkipper(int doc, long owningBucketOrd) throws IOException {
if (doc > skipper.maxDocID(0)) {
skipper.advance(doc);
}
upToSameBucket = false;

if (skipper.minDocID(0) > doc) {
// Corner case which happens if `doc` doesn't have a value and is between two
// intervals of
// the doc-value skip index.
upToInclusive = skipper.minDocID(0) - 1;
return;
}

upToInclusive = skipper.maxDocID(0);

// Now find the highest level where all docs map to the same bucket.
for (int level = 0; level < skipper.numLevels(); ++level) {
int totalDocsAtLevel = skipper.maxDocID(level) - skipper.minDocID(level) + 1;
long minBucket = preparedRounding.round(skipper.minValue(level));
long maxBucket = preparedRounding.round(skipper.maxValue(level));

if (skipper.docCount(level) == totalDocsAtLevel && minBucket == maxBucket) {
// All docs at this level have a value, and all values map to the same bucket.
upToInclusive = skipper.maxDocID(level);
upToSameBucket = true;
upToBucketIndex = bucketOrds.add(owningBucketOrd, maxBucket);
if (upToBucketIndex < 0) {
upToBucketIndex = -1 - upToBucketIndex;
}
} else {
break;
}
}
}

@Override
public void collect(int doc, long owningBucketOrd) throws IOException {
if (doc > upToInclusive) {
advanceSkipper(doc, owningBucketOrd);
}

if (upToSameBucket) {
aggregator.incrementBucketDocCount(upToBucketIndex, 1L);
sub.collect(doc, upToBucketIndex);
} else if (values.advanceExact(doc)) {
final long value = values.longValue();
long bucketIndex = bucketOrds.add(owningBucketOrd, preparedRounding.round(value));
if (bucketIndex < 0) {
bucketIndex = -1 - bucketIndex;
aggregator.collectExistingBucket(sub, doc, bucketIndex);
} else {
aggregator.collectBucket(sub, doc, bucketIndex);
}
}
}

@Override
public void collect(DocIdStream stream) throws IOException {
// This will only be called if its the top agg
collect(stream, 0);
}

@Override
public void collect(DocIdStream stream, long owningBucketOrd) throws IOException {
// This will only be called if its the sub aggregation
for (;;) {
int upToExclusive = upToInclusive + 1;
if (upToExclusive < 0) { // overflow
upToExclusive = Integer.MAX_VALUE;
}

if (upToSameBucket) {
if (sub == NO_OP_COLLECTOR) {
// stream.count maybe faster when we don't need to handle sub-aggs
long count = stream.count(upToExclusive);
aggregator.incrementBucketDocCount(upToBucketIndex, count);
} else {
int count = 0;
int[] docBuffer = new int[64];
int cnt = Integer.MAX_VALUE;
while (cnt != 0) {
cnt = stream.intoArray(upToExclusive, docBuffer);
sub.collect(docBuffer, upToBucketIndex);
Comment on lines +156 to +157

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think the issue here is when cnt is less than docBuffer.length, sub.collect(docBuffer, upToBucketIndex); will still iterate through entire docBuffer.

So we can either pass in size e.g.

sub.collect(docBuffer, cnt, upToBucketIndex);

Or create a new array, which I think will be sub optimal.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, I realized this yesterday. Missed adding a comment. Due to this issue advanceExact might be getting invoked for target > doc resulting in an error

count += cnt;
}
aggregator.incrementBucketDocCount(upToBucketIndex, count);
}
} else {
stream.forEach(upToExclusive, doc -> collect(doc, owningBucketOrd));
}

if (stream.mayHaveRemaining()) {
advanceSkipper(upToExclusive, owningBucketOrd);
} else {
break;
}
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
/*
* SPDX-License-Identifier: Apache-2.0
*
* The OpenSearch Contributors require contributions made to
* this file be licensed under the Apache-2.0 license or a
* compatible open source license.
*/

package org.opensearch.search.aggregations.bucket.filterrewrite;

import org.apache.lucene.index.LeafReaderContext;
import org.opensearch.search.aggregations.LeafBucketCollector;

import java.io.IOException;

/**
* Workaround for collectors that cannot handle DFS travel, i.e. changing owningBucketOrd)
*
* @opensearch.internal
*/
public interface BFSCollector {

LeafBucketCollector getBFSLeafCollector(LeafReaderContext ctx) throws IOException;
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
import org.apache.lucene.util.FixedBitSet;
import org.opensearch.search.aggregations.BucketCollector;
import org.opensearch.search.aggregations.LeafBucketCollector;
import org.opensearch.search.aggregations.bucket.filterrewrite.BFSCollector;
import org.opensearch.search.aggregations.bucket.filterrewrite.FilterRewriteOptimizationContext;
import org.opensearch.search.aggregations.bucket.filterrewrite.Ranges;

Expand Down Expand Up @@ -105,7 +106,12 @@ public void finalizePreviousRange() {
// trigger the sub agg collection for this range
try {
// build a new leaf collector for each bucket
LeafBucketCollector sub = collectableSubAggregators.getLeafCollector(leafCtx);
LeafBucketCollector sub = null;
if (collectableSubAggregators instanceof BFSCollector bfsCollector) {
sub = bfsCollector.getBFSLeafCollector(leafCtx);
} else {
sub = collectableSubAggregators.getLeafCollector(leafCtx);
}
sub.collect(DocIdStreamHelper.getDocIdStream(bitSet), bucketOrd);
logger.trace("collected sub aggregation for bucket {}", bucketOrd);
} catch (IOException e) {
Expand Down
Loading
Loading