This is an automated email from the ASF dual-hosted git repository.
xiangfu0 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 00e01874761 Keep tied ExprMin/ExprMax rows from every server when
merging serialized results (#19645)
00e01874761 is described below
commit 00e0187476192e622416084cf74f4b6d30fff338
Author: Xiang Fu <[email protected]>
AuthorDate: Fri Sep 25 20:29:13 2026 -0700
Keep tied ExprMin/ExprMax rows from every server when merging serialized
results (#19645)
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]