This is an automated email from the ASF dual-hosted git repository.

xiangfu0 pushed a commit to branch xiangfu0/agg-binding/00-exprminmax-tie-merge
in repository https://gitbox.apache.org/repos/asf/pinot.git

commit 1b75c064358e2cdcb3bef17ab9e412842eaff514
Author: Xiang Fu <[email protected]>
AuthorDate: Wed Sep 23 14:17:12 2026 -0700

    Keep tied ExprMin/ExprMax rows from every server when merging serialized 
results
    
    ExprMinMaxObject.merge switched a deserialized accumulator to mutable mode
    before copying its rows and key, so the copy saw zero rows. Tied rows from 
the
    first server were dropped and the merged key stayed null, which failed on 
the
    next tied merge or re-serialization. Copy the rows and key first, and honor 
the
    data block's null bitmap when reading immutable projections.
---
 .../utils/exprminmax/ExprMinMaxObject.java         |  17 ++-
 .../utils/exprminmax/ExprMinMaxObjectTest.java     |  92 ++++++++++++++
 .../org/apache/pinot/queries/ExprMinMaxTest.java   | 140 +++++++++------------
 3 files changed, 167 insertions(+), 82 deletions(-)

diff --git 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObject.java
 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObject.java
index e54257569f7..3f504ac0b51 100644
--- 
a/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObject.java
+++ 
b/pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObject.java
@@ -32,6 +32,7 @@ import org.apache.pinot.common.utils.DataSchema;
 import org.apache.pinot.core.common.datablock.DataBlockBuilder;
 import 
org.apache.pinot.core.query.aggregation.utils.ParentAggregationFunctionResultObject;
 import org.apache.pinot.segment.spi.memory.CompoundDataBuffer;
+import org.roaringbitmap.RoaringBitmap;
 
 
 @SuppressWarnings("rawtypes")
@@ -82,6 +83,8 @@ public class ExprMinMaxObject implements 
ParentAggregationFunctionResultObject {
   // used for ser/de
   private DataBlock _immutableMeasuringKeys;
   private DataBlock _immutableProjectionVals;
+  // Decoding a null bitmap copies it from the data block; cache once per 
column instead of once per projected cell.
+  private RoaringBitmap[] _immutableProjectionNullRows;
 
   public ExprMinMaxObject(DataSchema measuringSchema, DataSchema 
projectionSchema) {
     _isNull = true;
@@ -108,6 +111,10 @@ public class ExprMinMaxObject implements 
ParentAggregationFunctionResultObject {
 
     _sizeOfExtremumMeasuringKeys = _measuringSchema.size();
     _sizeOfExtremumProjectionVals = _projectionSchema.size();
+    _immutableProjectionNullRows = new 
RoaringBitmap[_sizeOfExtremumProjectionVals];
+    for (int i = 0; i < _sizeOfExtremumProjectionVals; i++) {
+      _immutableProjectionNullRows[i] = 
_immutableProjectionVals.getNullRowIds(i);
+    }
   }
 
   public static ExprMinMaxObject fromBytes(byte[] bytes)
@@ -250,6 +257,10 @@ public class ExprMinMaxObject implements 
ParentAggregationFunctionResultObject {
     if (_mutable) {
       return _extremumProjectionValues.get(rowId)[colId];
     } else {
+      RoaringBitmap nullRows = _immutableProjectionNullRows[colId];
+      if (nullRows != null && nullRows.contains(rowId)) {
+        return null;
+      }
       switch (_projectionSchema.getColumnDataType(colId)) {
         case INT:
           return _immutableProjectionVals.getInt(rowId, colId);
@@ -310,8 +321,8 @@ public class ExprMinMaxObject implements 
ParentAggregationFunctionResultObject {
       }
       // If the keys are equal, add the values of the other object to this 
object
       if (!_mutable) {
-        // If the result is immutable, we need to copy the values from the 
serialized result to the mutable result
-        _mutable = true;
+        // If the result is immutable, we need to copy the values from the 
serialized result to the mutable result.
+        // Read the rows and the key before switching to mutable, because both 
accessors depend on the mode.
         for (int i = 0; i < getNumberOfRows(); i++) {
           Object[] val = new Object[_sizeOfExtremumProjectionVals];
           for (int j = 0; j < _sizeOfExtremumProjectionVals; j++) {
@@ -319,6 +330,8 @@ public class ExprMinMaxObject implements 
ParentAggregationFunctionResultObject {
           }
           _extremumProjectionValues.add(val);
         }
+        _extremumMeasuringKeys = key;
+        _mutable = true;
       }
       for (int i = 0; i < other.getNumberOfRows(); i++) {
         Object[] val = new Object[_sizeOfExtremumProjectionVals];
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObjectTest.java
 
b/pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObjectTest.java
new file mode 100644
index 00000000000..7bf45e85c3d
--- /dev/null
+++ 
b/pinot-core/src/test/java/org/apache/pinot/core/query/aggregation/utils/exprminmax/ExprMinMaxObjectTest.java
@@ -0,0 +1,92 @@
+/**
+ * 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.utils.exprminmax;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.common.utils.DataSchema.ColumnDataType;
+import org.apache.pinot.core.operator.docvalsets.RowBasedBlockValSet;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+
+
+public class ExprMinMaxObjectTest {
+  private static final DataSchema MEASURING_SCHEMA =
+      new DataSchema(new String[]{"measure"}, new 
ColumnDataType[]{ColumnDataType.INT});
+  private static final DataSchema PROJECTION_SCHEMA =
+      new DataSchema(new String[]{"projection"}, new 
ColumnDataType[]{ColumnDataType.STRING});
+
+  @Test
+  public void testTiedSerializedResultsKeepAllRows()
+      throws IOException {
+    // The broker merges deserialized (immutable) server results. Tied rows 
from every server must survive, including
+    // the first one, and the merged result must still expose its key for the 
next merge.
+    ExprMinMaxObject merged = serialized(5, "a1", "a2")
+        .merge(serialized(5, "b1"), true)
+        .merge(serialized(5, "c1"), true)
+        .merge(serialized(3, "d1"), true);
+
+    assertEquals(merged.getExtremumKey(), new Comparable[]{5});
+    assertEquals(projections(merged), List.of("a1", "a2", "b1", "c1"));
+  }
+
+  @Test
+  public void testTiedSerializedResultSurvivesReserialization()
+      throws IOException {
+    ExprMinMaxObject merged = serialized(5, "a1").merge(serialized(5, "b1"), 
false);
+    ExprMinMaxObject roundTripped = 
ExprMinMaxObject.fromBytes(merged.toBytes());
+
+    assertEquals(roundTripped.getExtremumKey(), new Comparable[]{5});
+    assertEquals(projections(roundTripped), List.of("a1", "b1"));
+  }
+
+  /// Builds a single-server result through the segment accumulation path, 
then serializes it as a server would.
+  private static ExprMinMaxObject serialized(int key, String... projections)
+      throws IOException {
+    List<Object[]> rows = new ArrayList<>();
+    for (String projection : projections) {
+      rows.add(new Object[]{key, projection});
+    }
+    List<ExprMinMaxMeasuringValSetWrapper> measuring =
+        List.of(new ExprMinMaxMeasuringValSetWrapper(new 
RowBasedBlockValSet(ColumnDataType.INT, rows, 0, false)));
+    List<ExprMinMaxProjectionValSetWrapper> projection =
+        List.of(new ExprMinMaxProjectionValSetWrapper(new 
RowBasedBlockValSet(ColumnDataType.STRING, rows, 1, false)));
+    ExprMinMaxObject object = new ExprMinMaxObject(MEASURING_SCHEMA, 
PROJECTION_SCHEMA);
+    for (int i = 0; i < rows.size(); i++) {
+      int result = object.compareAndSetKey(measuring, i, true);
+      if (result > 0) {
+        object.setToNewVal(projection, i);
+      } else if (result == 0) {
+        object.addVal(projection, i);
+      }
+    }
+    return ExprMinMaxObject.fromBytes(object.toBytes());
+  }
+
+  private static List<Object> projections(ExprMinMaxObject object) {
+    List<Object> values = new ArrayList<>();
+    for (int i = 0; i < object.getNumberOfRows(); i++) {
+      values.add(object.getField(i, 0));
+    }
+    return values;
+  }
+}
diff --git 
a/pinot-core/src/test/java/org/apache/pinot/queries/ExprMinMaxTest.java 
b/pinot-core/src/test/java/org/apache/pinot/queries/ExprMinMaxTest.java
index 82a96a5359e..756bab94330 100644
--- a/pinot-core/src/test/java/org/apache/pinot/queries/ExprMinMaxTest.java
+++ b/pinot-core/src/test/java/org/apache/pinot/queries/ExprMinMaxTest.java
@@ -58,7 +58,8 @@ import static org.testng.Assert.assertTrue;
 import static org.testng.Assert.fail;
 
 
-/// Queries test for exprmin/exprmax functions.
+/// Queries test for exprmin/exprmax functions. BaseQueriesTest queries two 
segments on each of two servers, so
+/// each qualifying source row occurs four times. Tied extrema must retain all 
four copies after serialized merging.
 public class ExprMinMaxTest extends BaseQueriesTest {
   private static final File INDEX_DIR = new File(FileUtils.getTempDirectory(), 
"ExprMinMaxTest");
   private static final String RAW_TABLE_NAME = "testTable";
@@ -217,7 +218,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
 
     assertEquals(rows.get(0)[0], 999L);
     assertEquals(rows.get(1)[0], 999L);
-    assertEquals(rows.size(), 2);
+    assertEquals(rows.size(), 4);
 
     // Inter segment data type test
     query = "SELECT expr_max(longColumn, intColumn), expr_max(floatColumn, 
intColumn), "
@@ -244,7 +245,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     assertEquals(resultTable.getDataSchema().getColumnName(9), 
"exprmax(mvDoubleColumn,bigDecimalColumn)");
     assertEquals(resultTable.getDataSchema().getColumnName(10), 
"exprmax(jsonColumn,bigDecimalColumn)");
 
-    assertEquals(rows.size(), 2);
+    assertEquals(rows.size(), 4);
     assertEquals(rows.get(0)[0], 999L);
     assertEquals(rows.get(1)[0], 999L);
     assertEquals(rows.get(0)[1], 999.5F);
@@ -275,8 +276,8 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 2000);
-    for (int i = 0; i < 2000; i++) {
+    assertEquals(rows.size(), 4000);
+    for (int i = 0; i < 4000; i++) {
       assertEquals(rows.get(i)[0], 1);
     }
 
@@ -289,27 +290,14 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 4);
-
-    assertEquals(rows.get(0)[0], 7996000D);
-    assertEquals(rows.get(0)[1], 8D);
-    assertEquals(rows.get(0)[2], "a11");
-    assertEquals(rows.get(0)[3], 8D);
+    assertEquals(rows.size(), 8);
 
-    assertEquals(rows.get(1)[0], 7996000D);
-    assertEquals(rows.get(1)[1], 18D);
-    assertEquals(rows.get(1)[2], "a11");
-    assertEquals(rows.get(1)[3], 8D);
-
-    assertEquals(rows.get(2)[0], 7996000D);
-    assertEquals(rows.get(2)[1], 8D);
-    assertEquals(rows.get(2)[2], "a11");
-    assertNull(rows.get(2)[3]);
-
-    assertEquals(rows.get(3)[0], 7996000D);
-    assertEquals(rows.get(3)[1], 18D);
-    assertEquals(rows.get(3)[2], "a11");
-    assertNull(rows.get(3)[3]);
+    for (int i = 0; i < 8; i++) {
+      assertEquals(rows.get(i)[0], 7996000D);
+      assertEquals(rows.get(i)[1], i % 2 == 0 ? 8D : 18D);
+      assertEquals(rows.get(i)[2], "a11");
+      assertEquals(rows.get(i)[3], i < 4 ? 8D : null);
+    }
 
     // Test transformation function inside exprmax/exprmin, for both 
projection and measuring
     // the max of 3000x-x^2 is 2250000, which is the max of 3000x-x^2
@@ -323,22 +311,14 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 4);
+    assertEquals(rows.size(), 8);
 
-    assertEquals(rows.get(0)[0], 7996000D);
-    assertEquals(rows.get(0)[1], 1500D);
-    assertEquals(rows.get(0)[2], 2250000D);
-    assertEquals(rows.get(0)[3], "bb11");
-    assertEquals(rows.get(1)[0], 7996000D);
-    assertEquals(rows.get(1)[1], 1500D);
-    assertEquals(rows.get(1)[2], 2250000D);
-    assertEquals(rows.get(1)[3], "bb11");
-    assertEquals(rows.get(2)[0], 7996000D);
-    assertNull(rows.get(2)[1]);
-    assertEquals(rows.get(2)[3], "bb11");
-    assertEquals(rows.get(3)[0], 7996000D);
-    assertNull(rows.get(3)[1]);
-    assertEquals(rows.get(3)[3], "bb11");
+    for (int i = 0; i < 8; i++) {
+      assertEquals(rows.get(i)[0], 7996000D);
+      assertEquals(rows.get(i)[1], i < 4 ? 1500D : null);
+      assertEquals(rows.get(i)[2], i < 4 ? 2250000D : null);
+      assertEquals(rows.get(i)[3], "bb11");
+    }
 
     // Inter segment mix aggregation function with CASE statement
     query = "SELECT exprmin(stringColumn, CASE WHEN stringColumn = 'a33' THEN 
'b' WHEN stringColumn = 'a22' THEN 'a' "
@@ -350,9 +330,9 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 4);
+    assertEquals(rows.size(), 8);
 
-    for (int i = 0; i < 4; i++) {
+    for (int i = 0; i < 8; i++) {
       assertEquals(rows.get(i)[0], "a22");
       assertEquals(rows.get(i)[1], "a");
     }
@@ -363,7 +343,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     brokerResponse = getBrokerResponse(query);
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
-    assertEquals(rows.size(), 2);
+    assertEquals(rows.size(), 4);
     assertEquals(rows.get(0)[0], "30");
     assertEquals(rows.get(1)[0], "30");
 
@@ -373,7 +353,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     brokerResponse = getBrokerResponse(query);
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
-    assertEquals(rows.size(), 2);
+    assertEquals(rows.size(), 4);
     String[] expectedMvBytes = {"30", "31", "32"};
     assertEquals(rows.get(0)[0], expectedMvBytes);
     assertEquals(rows.get(1)[0], expectedMvBytes);
@@ -389,7 +369,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     ResultTable resultTable = brokerResponse.getResultTable();
     List<Object[]> rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 4);
+    assertEquals(rows.size(), 8);
 
     assertEquals(rows.get(0)[0], 0);
     assertEquals(rows.get(1)[0], 1200);
@@ -404,7 +384,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 2);
+    assertEquals(rows.size(), 4);
 
     assertEquals(rows.get(0)[0], 0);
     assertEquals(rows.get(1)[0], 0);
@@ -418,7 +398,7 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 2);
+    assertEquals(rows.size(), 4);
 
     assertEquals(rows.get(0)[0], 1200);
     assertEquals(rows.get(1)[0], 1200);
@@ -453,10 +433,10 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     ResultTable resultTable = brokerResponse.getResultTable();
     List<Object[]> rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 10);
+    assertEquals(rows.size(), 20);
 
-    for (int i = 0; i < 10; i++) {
-      int group = ((i + 2) / 2) % 5;
+    for (int i = 0; i < 20; i++) {
+      int group = (i / 4 + 1) % 5;
       assertEquals(rows.get(i)[0], group);
       assertEquals(rows.get(i)[1], 995L + group);
     }
@@ -470,19 +450,18 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 24);
+    assertEquals(rows.size(), 48);
 
-    for (int i = 0; i < 22; i++) {
-      double group = Math.pow(2, i / 2);
+    for (int i = 0; i < 44; i++) {
+      double group = Math.pow(2, i / 4);
       assertEquals(rows.get(i)[0], (int) group);
       assertEquals(rows.get(i)[1], group - 1);
     }
 
-    assertEquals(rows.get(22)[0], 2048);
-    assertEquals(rows.get(22)[1], 1999D);
-
-    assertEquals(rows.get(23)[0], 2048);
-    assertEquals(rows.get(23)[1], 1999D);
+    for (int i = 44; i < 48; i++) {
+      assertEquals(rows.get(i)[0], 2048);
+      assertEquals(rows.get(i)[1], 1999D);
+    }
 
     // MV inter segment group by
     query = "SELECT groupByMVIntColumn, expr_min(doubleColumn, intColumn) FROM 
testTable GROUP BY groupByMVIntColumn";
@@ -491,19 +470,18 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
 
-    assertEquals(rows.size(), 20);
+    assertEquals(rows.size(), 40);
 
-    for (int i = 0; i < 18; i++) {
-      int group = i / 2 + 1;
+    for (int i = 0; i < 36; i++) {
+      int group = i / 4 + 1;
       assertEquals(rows.get(i)[0], group);
       assertEquals(rows.get(i)[1], (double) group - 1);
     }
 
-    assertEquals(rows.get(18)[0], 0);
-    assertEquals(rows.get(18)[1], 0D);
-
-    assertEquals(rows.get(19)[0], 0);
-    assertEquals(rows.get(19)[1], 0D);
+    for (int i = 36; i < 40; i++) {
+      assertEquals(rows.get(i)[0], 0);
+      assertEquals(rows.get(i)[1], 0D);
+    }
 
     // MV inter segment group by with projection on MV column
     query = "SELECT groupByMVIntColumn, expr_min(mvIntColumn, intColumn), "
@@ -512,18 +490,20 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     brokerResponse = getBrokerResponse(query);
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
-    assertEquals(rows.size(), 20);
+    assertEquals(rows.size(), 40);
 
-    for (int i = 0; i < 18; i++) {
-      int group = i / 2 + 1;
+    for (int i = 0; i < 36; i++) {
+      int group = i / 4 + 1;
       assertEquals(rows.get(i)[0], group);
       assertEquals(rows.get(i)[1], new Object[]{group - 1, group, group + 1});
       assertEquals(rows.get(i)[2], new Object[]{"a199" + group, "a199" + group 
+ 1, "a199" + group + 2});
     }
 
-    assertEquals(rows.get(18)[0], 0);
-    assertEquals(rows.get(18)[1], new Object[]{0, 1, 2});
-    assertEquals(rows.get(18)[2], new Object[]{"a1999", "a19991", "a19992"});
+    for (int i = 36; i < 40; i++) {
+      assertEquals(rows.get(i)[0], 0);
+      assertEquals(rows.get(i)[1], new Object[]{0, 1, 2});
+      assertEquals(rows.get(i)[2], new Object[]{"a1999", "a19991", "a19992"});
+    }
   }
 
   @Test
@@ -538,10 +518,10 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     BrokerResponse brokerResponse = getBrokerResponse(query);
     ResultTable resultTable = brokerResponse.getResultTable();
     List<Object[]> rows = resultTable.getRows();
-    assertEquals(rows.size(), 14);
-    assertEquals(rows.get(4)[0], "a33");
-    assertEquals(rows.get(4)[1], new Object[]{20, 21, 22});
-    assertEquals(rows.get(4)[2], new Object[]{27});
+    assertEquals(rows.size(), 28);
+    assertEquals(rows.get(8)[0], "a33");
+    assertEquals(rows.get(8)[1], new Object[]{20, 21, 22});
+    assertEquals(rows.get(8)[2], new Object[]{27});
 
     //  TODO: The following query works because whenever we find an empty 
array in the result, we use null
     //        (see exprminMaxProjectionValSetWrapper). Ideally, we should be 
able to serialize empty array.
@@ -554,10 +534,10 @@ public class ExprMinMaxTest extends BaseQueriesTest {
     brokerResponse = getBrokerResponse(query);
     resultTable = brokerResponse.getResultTable();
     rows = resultTable.getRows();
-    assertEquals(rows.size(), 20);
-    assertEquals(rows.get(8)[0], "a33");
-    assertEquals(rows.get(8)[1], new Object[]{20, 21, 22});
-    assertEquals(rows.get(8)[2], new Object[]{});
+    assertEquals(rows.size(), 40);
+    assertEquals(rows.get(16)[0], "a33");
+    assertEquals(rows.get(16)[1], new Object[]{20, 21, 22});
+    assertNull(rows.get(16)[2]);
   }
 
   @Test


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

Reply via email to