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]
