xiangfu0 commented on code in PR #19587:
URL: https://github.com/apache/pinot/pull/19587#discussion_r4212836522


##########
pinot-common/src/main/java/org/apache/pinot/common/utils/MutableRoaringBitmapUnion.java:
##########
@@ -0,0 +1,207 @@
+/**
+ * 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.pinot.common.utils;
+
+import java.util.Objects;
+import org.roaringbitmap.buffer.ImmutableRoaringBitmap;
+import org.roaringbitmap.buffer.MappeableContainerPointer;
+import org.roaringbitmap.buffer.MutableRoaringBitmap;
+
+
+/// Interim stand-in for `org.roaringbitmap.buffer.MutableRoaringBitmapUnion`, 
the buffer counterpart of
+/// [RoaringBitmapUnion]: same name and methods as the class proposed to 
RoaringBitmap, so that the swap is an import
+/// change plus deleting this class. It is internal to Pinot and goes away 
with the swap. See [RoaringBitmapUnion] for
+/// the contract and the swap steps.
+///
+/// Inputs may be read-only bitmaps over memory-mapped buffers; they are only 
read.
+///
+/// Differences from the library class:
+///
+/// - The accumulated state is a private subclass of [MutableRoaringBitmap] 
that reaches the protected lazy primitives
+///   through inheritance; [#get()] and [#take()] hand out instances of it.
+/// - [#takeOwnership(MutableRoaringBitmap)] adopts without copying only a 
bitmap that a union handed out; any other
+///   bitmap is copied. Callers should still treat the argument as 
relinquished, and should not pass subclasses.
+/// - [#add(int)] repairs pending lazy state before inserting and does not 
re-encode a run container that received
+///   the value.
+/// - An input is unioned eagerly while the accumulator holds few values per 
container, or when it would insert many
+///   container keys that are new to the accumulator, because the released 
lazy union only pays off for dense
+///   containers and inserts new keys one at a time.
+///
+/// Instances are not thread-safe.
+public final class MutableRoaringBitmapUnion {
+  // See RoaringBitmapUnion
+  private static final int MIN_VALUES_PER_CONTAINER_FOR_LAZY_UNION = 1024;
+  private static final int MAX_NEW_KEYS_FOR_LAZY_UNION = 4;
+
+  // The accumulated bitmap. Never null.
+  private LazyBitmap _bitmap;
+  // Whether the bitmap holds lazy state that must be repaired before it is 
read
+  private boolean _dirty;
+  // Whether an alias returned by get() is outstanding and must not be mutated
+  private boolean _published;
+  // Number of values added so far, counting duplicates: an upper bound of the 
cardinality that is known without
+  // repairing, used to estimate how dense the containers are
+  private long _numValuesAdded;
+  // A cardinality observed on repaired state. Unions never remove values, so 
this remains a lower bound and proves
+  // density without repairing every lazy union. New container keys can 
invalidate that density proof.
+  private long _lastKnownCardinality;
+
+  /// Creates an empty union.
+  public MutableRoaringBitmapUnion() {
+    _bitmap = new LazyBitmap();
+  }
+
+  private MutableRoaringBitmapUnion(LazyBitmap adopted) {
+    _bitmap = adopted;
+    // Adopted state is normalized on the first read, like the library class 
does
+    _dirty = true;
+    _numValuesAdded = adopted.getLongCardinality();
+    _lastKnownCardinality = _numValuesAdded;
+  }
+
+  /// Creates a union whose initial state is the given bitmap. The caller 
relinquishes the instance: it must not be
+  /// used again except through the union.
+  public static MutableRoaringBitmapUnion takeOwnership(MutableRoaringBitmap 
bitmap) {

Review Comment:
   Removed in 
[e6fc4ab](https://github.com/apache/pinot/commit/e6fc4abe27b538afea2fb35f5e194849960f3003):
 both `takeOwnership` factories, the unused buffered adoption constructor, and 
`MutableRoaringBitmapUnion.add(int)`. The heap `add(int)` and direct 
deserialization path remain because production callers use them. The docs now 
describe only the compatible API subset Pinot needs.
   
   The existing copy-on-write, take/reset, sparse/dense and skewed-density 
regressions remain. Container-identity assertions use a narrow test-only 
inspection instead of retaining an unused production entry point; both adapted 
skew tests still fail on the prior implementation and pass with the fix. All 31 
focused tests and scoped quality checks pass; independent review found no 
issues. Hosted CI is running for this commit.
   



##########
pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/function/DistinctCountBitmapLazyUnionTest.java:
##########
@@ -0,0 +1,329 @@
+/**
+ * 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.pinot.core.query.aggregation.function;
+
+import java.nio.ByteBuffer;
+import java.util.Map;
+import java.util.Random;
+import javax.annotation.Nullable;
+import org.apache.pinot.common.CustomObject;
+import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.utils.RoaringBitmapUtils;
+import org.apache.pinot.core.common.BlockValSet;
+import org.apache.pinot.core.query.aggregation.AggregationResultHolder;
+import 
org.apache.pinot.core.query.aggregation.function.AggregationFunction.SerializedIntermediateResult;
+import org.apache.pinot.core.query.aggregation.groupby.GroupByResultHolder;
+import org.apache.pinot.spi.data.FieldSpec.DataType;
+import org.mockito.Mockito;
+import org.roaringbitmap.RoaringBitmap;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertSame;
+import static org.testng.Assert.assertTrue;
+
+
+/// Covers the serialized-bitmap (`BYTES`) aggregation paths of 
[DistinctCountBitmapAggregationFunction], which union
+/// input bitmaps lazily and repair the accumulator at extraction. The input 
cardinalities are chosen to cross the
+/// container thresholds of the lazy union (array containers promote to bitmap 
containers past 1024 combined
+/// cardinality; repair converts back to array containers at up to 4096), so 
both promoted and non-promoted
+/// accumulator states are verified against an eagerly unioned reference. 
Intermediate-result merges must remain
+/// valid for cardinality reads and serialization immediately after every 
merge.
+public class DistinctCountBitmapLazyUnionTest {
+  private static final ExpressionContext EXPRESSION = 
ExpressionContext.forIdentifier("bitmapCol");
+  private static final Random RANDOM = new Random(42);
+
+  private static byte[][] serializedBitmaps(int numBitmaps, int 
valuesPerBitmap, int maxValue) {
+    byte[][] serialized = new byte[numBitmaps][];
+    for (int i = 0; i < numBitmaps; i++) {
+      RoaringBitmap bitmap = new RoaringBitmap();
+      for (int j = 0; j < valuesPerBitmap; j++) {
+        bitmap.add(RANDOM.nextInt(maxValue));
+      }
+      serialized[i] = RoaringBitmapUtils.serialize(bitmap);
+    }
+    return serialized;
+  }
+
+  private static RoaringBitmap eagerUnion(byte[][] serialized, int from, int 
to) {
+    RoaringBitmap expected = new RoaringBitmap();
+    for (int i = from; i < to; i++) {
+      expected.or(RoaringBitmapUtils.deserialize(serialized[i]));
+    }
+    return expected;
+  }
+
+  private static Map<ExpressionContext, BlockValSet> 
mockBlockValSetMap(byte[][] serialized) {
+    return mockBlockValSetMap(serialized, null);
+  }
+
+  private static Map<ExpressionContext, BlockValSet> 
mockBlockValSetMap(byte[][] serialized,
+      @Nullable RoaringBitmap nullBitmap) {
+    BlockValSet blockValSet = Mockito.mock(BlockValSet.class);
+    Mockito.when(blockValSet.getValueType()).thenReturn(DataType.BYTES);
+    Mockito.when(blockValSet.isSingleValue()).thenReturn(true);
+    Mockito.when(blockValSet.getBytesValuesSV()).thenReturn(serialized);
+    Mockito.when(blockValSet.getNullBitmap()).thenReturn(nullBitmap);
+    return Map.of(EXPRESSION, blockValSet);
+  }
+
+  @Test
+  public void testAggregateSkipsNullRowsAcrossRanges() {
+    // With null handling enabled the block is consumed as several non-null 
ranges; the accumulator is re-read from
+    // the holder at the start of each range and written back at its end, and 
null rows never reach the union
+    byte[][] serialized = serializedBitmaps(120, 50, 100_000);
+    RoaringBitmap nullBitmap = RoaringBitmap.bitmapOf(0, 1, 17, 40, 41, 42, 
99, 119);
+    DistinctCountBitmapAggregationFunction function =
+        new DistinctCountBitmapAggregationFunction(EXPRESSION, true);
+    AggregationResultHolder holder = function.createAggregationResultHolder();
+    int blockSize = 60;
+    for (int from = 0; from < serialized.length; from += blockSize) {
+      byte[][] block = new byte[blockSize][];
+      System.arraycopy(serialized, from, block, 0, blockSize);
+      RoaringBitmap blockNulls = new RoaringBitmap();
+      for (int i = 0; i < blockSize; i++) {
+        if (nullBitmap.contains(from + i)) {
+          blockNulls.add(i);
+        }
+      }
+      function.aggregate(blockSize, holder, mockBlockValSetMap(block, 
blockNulls));
+    }
+
+    RoaringBitmap expected = new RoaringBitmap();
+    for (int i = 0; i < serialized.length; i++) {
+      if (!nullBitmap.contains(i)) {
+        expected.or(RoaringBitmapUtils.deserialize(serialized[i]));
+      }
+    }
+    RoaringBitmap result = function.extractAggregationResult(holder);
+    assertEquals(result, expected);
+    assertEquals(result.getCardinality(), expected.getCardinality());
+    assertFalse(result.isEmpty());
+  }
+
+  @Test
+  public void testAggregateAcrossBlocks() {
+    // Enough overlap-heavy inputs to promote accumulator containers to (lazy) 
bitmap containers, split into
+    // multiple aggregate() calls to verify the accumulator stays valid across 
blocks until extraction. An early
+    // extraction after the first block guards the published-bitmap safety 
net: no operator extracts and then keeps
+    // aggregating today, but if one did, the published bitmap must stay 
immutable and later extraction must still
+    // be complete.
+    byte[][] serialized = serializedBitmaps(200, 100, 200_000);
+    DistinctCountBitmapAggregationFunction function =
+        new DistinctCountBitmapAggregationFunction(EXPRESSION, false);
+    AggregationResultHolder holder = function.createAggregationResultHolder();
+    int blockSize = 50;
+    RoaringBitmap published = null;
+    RoaringBitmap publishedCopy = null;
+    for (int from = 0; from < serialized.length; from += blockSize) {
+      byte[][] block = new byte[blockSize][];
+      System.arraycopy(serialized, from, block, 0, blockSize);
+      function.aggregate(blockSize, holder, mockBlockValSetMap(block));
+      if (published != null) {
+        assertEquals(published, publishedCopy);
+        assertEquals(published.getCardinality(), 
publishedCopy.getCardinality());
+      }
+      if (from == 0) {
+        published = function.extractAggregationResult(holder);
+        assertEquals(published, eagerUnion(serialized, 0, blockSize));
+        // Repeated extraction returns the same published instance
+        assertSame(function.extractAggregationResult(holder), published);
+        publishedCopy = published.clone();
+      }
+    }
+
+    RoaringBitmap result = function.extractAggregationResult(holder);
+    RoaringBitmap expected = eagerUnion(serialized, 0, serialized.length);
+    assertEquals(result, expected);
+    assertEquals(result.getCardinality(), expected.getCardinality());
+    // The repaired accumulator must serialize into a form that round-trips
+    
assertEquals(RoaringBitmapUtils.deserialize(RoaringBitmapUtils.serialize(result)),
 expected);
+  }
+
+  @Test
+  public void testAggregateGroupBySV() {
+    int numGroups = 4;
+    byte[][] serialized = serializedBitmaps(400, 50, 100_000);
+    int[] groupKeys = new int[serialized.length];
+    for (int i = 0; i < groupKeys.length; i++) {
+      groupKeys[i] = i % numGroups;
+    }
+
+    DistinctCountBitmapAggregationFunction function =
+        new DistinctCountBitmapAggregationFunction(EXPRESSION, false);
+    GroupByResultHolder holder = function.createGroupByResultHolder(numGroups, 
numGroups);
+    function.aggregateGroupBySV(serialized.length, groupKeys, holder, 
mockBlockValSetMap(serialized));
+
+    for (int groupKey = 0; groupKey < numGroups; groupKey++) {
+      RoaringBitmap expected = new RoaringBitmap();
+      for (int i = groupKey; i < serialized.length; i += numGroups) {
+        expected.or(RoaringBitmapUtils.deserialize(serialized[i]));
+      }
+      RoaringBitmap result = function.extractGroupByResult(holder, groupKey);
+      assertEquals(result, expected, "group " + groupKey);
+      assertEquals(result.getCardinality(), expected.getCardinality(), "group 
" + groupKey);
+    }
+  }
+
+  @Test
+  public void testAggregateGroupByMV() {
+    // Every row belongs to multiple groups, so the same deserialized input 
bitmap is unioned into one group's
+    // accumulator and cloned as another group's initial accumulator — 
verifies no aliasing between accumulators
+    int numGroups = 3;
+    byte[][] serialized = serializedBitmaps(150, 80, 150_000);
+    int[][] groupKeys = new int[serialized.length][];
+    for (int i = 0; i < groupKeys.length; i++) {
+      groupKeys[i] = new int[]{i % numGroups, (i + 1) % numGroups};
+    }
+
+    DistinctCountBitmapAggregationFunction function =
+        new DistinctCountBitmapAggregationFunction(EXPRESSION, false);
+    GroupByResultHolder holder = function.createGroupByResultHolder(numGroups, 
numGroups);
+    function.aggregateGroupByMV(serialized.length, groupKeys, holder, 
mockBlockValSetMap(serialized));
+
+    for (int groupKey = 0; groupKey < numGroups; groupKey++) {
+      RoaringBitmap expected = new RoaringBitmap();
+      for (int i = 0; i < serialized.length; i++) {
+        if (i % numGroups == groupKey || (i + 1) % numGroups == groupKey) {
+          expected.or(RoaringBitmapUtils.deserialize(serialized[i]));
+        }
+      }
+      RoaringBitmap result = function.extractGroupByResult(holder, groupKey);
+      assertEquals(result, expected, "group " + groupKey);
+      assertEquals(result.getCardinality(), expected.getCardinality(), "group 
" + groupKey);
+    }
+  }
+
+  @Test
+  public void testSmallBitmapsStayCorrectBelowPromotionThreshold() {
+    // All containers stay small array containers: the lazy path must degrade 
to plain array unions
+    byte[][] serialized = serializedBitmaps(50, 5, 1_000);
+    DistinctCountBitmapAggregationFunction function =
+        new DistinctCountBitmapAggregationFunction(EXPRESSION, false);
+    AggregationResultHolder holder = function.createAggregationResultHolder();
+    function.aggregate(serialized.length, holder, 
mockBlockValSetMap(serialized));
+
+    RoaringBitmap expected = eagerUnion(serialized, 0, serialized.length);
+    assertEquals(function.extractAggregationResult(holder), expected);
+  }
+
+  @Test
+  public void testMergeSparseBitmapsAcrossUnsignedRange() {

Review Comment:
   Removed in 
[e6fc4ab](https://github.com/apache/pinot/commit/e6fc4abe27b538afea2fb35f5e194849960f3003).
 I dropped the three merge-only cases from `DistinctCountBitmapLazyUnionTest` 
and restored `IndexedTableTest` and `AggregateOperatorTest` to their upstream 
contents, removing these unrelated additions from the PR. `merge()` remains 
unchanged.
   
   The five remaining BYTES aggregation tests cover null ranges, multi-block 
accumulation and stable extraction, SV/MV grouping with input reuse, small 
bitmaps and serialization. Those tests, both union suites and the StarTree 
bitmap value-aggregation suite pass: 31 tests, zero failures/errors/skips. 
Scoped quality checks and independent review also passed.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to