This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang 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 4f495df35c8 Fix wrong results for multi-key GROUP BY with transform
expressions (#19663)
4f495df35c8 is described below
commit 4f495df35c83f4446418cd45da79f10812c91799
Author: zeronerdzerogeekzerocool <[email protected]>
AuthorDate: Sat Oct 3 14:12:52 2026 -0700
Fix wrong results for multi-key GROUP BY with transform expressions (#19663)
---
.../core/query/utils/OrderByComparatorFactory.java | 17 +++--
.../combine/SortedGroupByCombineOperatorsTest.java | 89 ++++++++++++++++++++++
.../query/utils/OrderByComparatorFactoryTest.java | 38 +++++++++
3 files changed, 137 insertions(+), 7 deletions(-)
diff --git
a/pinot-core/src/main/java/org/apache/pinot/core/query/utils/OrderByComparatorFactory.java
b/pinot-core/src/main/java/org/apache/pinot/core/query/utils/OrderByComparatorFactory.java
index 6c044df5606..6368181f6c3 100644
---
a/pinot-core/src/main/java/org/apache/pinot/core/query/utils/OrderByComparatorFactory.java
+++
b/pinot-core/src/main/java/org/apache/pinot/core/query/utils/OrderByComparatorFactory.java
@@ -70,11 +70,15 @@ public class OrderByComparatorFactory {
return (k1, k2) -> valueComparator.compare(k1.getValues(), k2.getValues());
}
- private static Map<String, Integer>
getGroupByExpressionIndexMap(List<ExpressionContext> groupByExpressions) {
- Map<String, Integer> groupByExpressionIndexMap = new HashMap<>();
+ private static Map<ExpressionContext, Integer> getGroupByExpressionIndexMap(
+ List<ExpressionContext> groupByExpressions) {
+ Map<ExpressionContext, Integer> groupByExpressionIndexMap = new
HashMap<>();
int numGroupByExpressions = groupByExpressions.size();
for (int i = 0; i < numGroupByExpressions; i++) {
- groupByExpressionIndexMap.put(groupByExpressions.get(i).getIdentifier(),
i);
+ // NOTE: Key on the whole expression, not on getIdentifier().
getIdentifier() is null for anything that is not a
+ // plain column reference, so keying on it collapses every
transform group-by key onto a single null entry
+ // and makes all ORDER BY expressions resolve to the same column
index.
+ groupByExpressionIndexMap.put(groupByExpressions.get(i), i);
}
return groupByExpressionIndexMap;
}
@@ -93,12 +97,11 @@ public class OrderByComparatorFactory {
/// Add an index for each orderby expression with respect to its position in
the group keys
private static List<OrderByExpressionWithIndex>
getGroupKeyOrderByExpressionFromRowOrderByExpressions(
List<OrderByExpressionContext> rowOrderByExpressions,
List<ExpressionContext> groupByExpressions) {
- Map<String, Integer> groupByExpressionIndexMap =
getGroupByExpressionIndexMap(groupByExpressions);
+ Map<ExpressionContext, Integer> groupByExpressionIndexMap =
getGroupByExpressionIndexMap(groupByExpressions);
List<OrderByExpressionWithIndex> result = new ArrayList<>();
// get index wrt group key for each order by expression
- rowOrderByExpressions.forEach(expr ->
- result.add(
- new OrderByExpressionWithIndex(expr,
groupByExpressionIndexMap.get(expr.getExpression().getIdentifier()))));
+ rowOrderByExpressions.forEach(
+ expr -> result.add(new OrderByExpressionWithIndex(expr,
groupByExpressionIndexMap.get(expr.getExpression()))));
return result;
}
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/SortedGroupByCombineOperatorsTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/SortedGroupByCombineOperatorsTest.java
index da7114eaf57..3591b57be8e 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/SortedGroupByCombineOperatorsTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/operator/combine/SortedGroupByCombineOperatorsTest.java
@@ -22,6 +22,7 @@ import com.google.common.base.Preconditions;
import java.io.File;
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -73,6 +74,16 @@ public class SortedGroupByCombineOperatorsTest {
private static final Schema SCHEMA =
new Schema.SchemaBuilder().addSingleValueDimension(INT_COLUMN,
FieldSpec.DataType.INT).build();
+ // Schema/table for testSafeTrim*CombineWithTransformGroupByKeys below: two
independent key columns so a
+ // multi-key transform GROUP BY can be exercised.
+ private static final String COL_A = "colA";
+ private static final String COL_B = "colB";
+ private static final TableConfig TRANSFORM_KEY_TABLE_CONFIG =
+ new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME).build();
+ private static final Schema TRANSFORM_KEY_SCHEMA =
+ new Schema.SchemaBuilder().addSingleValueDimension(COL_A,
FieldSpec.DataType.INT)
+ .addSingleValueDimension(COL_B, FieldSpec.DataType.INT).build();
+
private static final PlanMaker PLAN_MAKER = new InstancePlanMakerImplV2();
private static final ExecutorService EXECUTOR =
Executors.newFixedThreadPool(4);
@@ -286,6 +297,84 @@ public class SortedGroupByCombineOperatorsTest {
}
}
+ // ----
+ // Multi-key transform GROUP BY (regression test for the
OrderByComparatorFactory bug where the ORDER BY ->
+ // GROUP BY index map was keyed by ExpressionContext#getIdentifier(), which
is null for any transform expression.
+ // With two transform-based GROUP BY keys, both collapsed onto the same null
map entry, so the merge comparator
+ // resolved every ORDER BY expression to the same column index instead of
its own.)
+ @Test
+ public void testSafeTrimPairWiseCombineWithTransformGroupByKeys()
+ throws Exception {
+ assertTransformGroupByKeysCombine(1);
+ }
+
+ @Test
+ public void testSafeTrimSequentialCombineWithTransformGroupByKeys()
+ throws Exception {
+ assertTransformGroupByKeysCombine(10_000_000);
+ }
+
+ /// segment 1 has 5 rows with (colA=0, colB=0); segment 2 has 5 rows with
(colA=1, colB=0). Both GROUP BY keys are
+ /// transform expressions, so both have a null identifier. Before the fix,
that made the merge comparator compare
+ /// every record only by its second key (colB), which is 0 in both segments,
so `SortedRecordsMerger` treated the
+ /// two distinct (colA, colB) groups as equal and merged their counts into a
single row instead of keeping two.
+ private void assertTransformGroupByKeysCombine(int
sortAggregateSingleThreadedNumSegmentsThreshold)
+ throws Exception {
+ IndexSegment segment1 = createTransformKeySegment("transformKeySegment_0",
0, 0, 5);
+ IndexSegment segment2 = createTransformKeySegment("transformKeySegment_1",
1, 0, 5);
+ try {
+ String query = "SET sortAggregateSingleThreadedNumSegmentsThreshold="
+ + sortAggregateSingleThreadedNumSegmentsThreshold + "; "
+ + "SELECT CEIL(colA), CEIL(colB), COUNT(*) FROM testTable GROUP BY
CEIL(colA), CEIL(colB) "
+ + "ORDER BY CEIL(colA), CEIL(colB) LIMIT 100";
+ QueryContext queryContext =
QueryContextConverterUtils.getQueryContext(query);
+ List<PlanNode> planNodes = new ArrayList<>(2);
+ planNodes.add(PLAN_MAKER.makeSegmentPlanNode(new
SegmentContext(segment1), queryContext));
+ planNodes.add(PLAN_MAKER.makeSegmentPlanNode(new
SegmentContext(segment2), queryContext));
+ queryContext.setEndTimeMs(
+ System.currentTimeMillis() +
CommonConstants.Server.DEFAULT_QUERY_EXECUTOR_TIMEOUT_MS);
+ CombinePlanNode combinePlanNode = new CombinePlanNode(planNodes,
queryContext, EXECUTOR, null);
+ BaseCombineOperator combineOperator = combinePlanNode.run();
+ GroupByResultsBlock combineResult = (GroupByResultsBlock)
combineOperator.nextBlock();
+
+ List<Object[]> rows = combineResult.getRows();
+ assertEquals(rows.size(), 2, "expected 2 distinct (colA, colB) groups,
got: " + Arrays.deepToString(
+ rows.toArray()));
+ assertEquals(rows.get(0)[0], 0.0);
+ assertEquals(rows.get(0)[1], 0.0);
+ assertEquals(rows.get(0)[2], 5L);
+ assertEquals(rows.get(1)[0], 1.0);
+ assertEquals(rows.get(1)[1], 0.0);
+ assertEquals(rows.get(1)[2], 5L);
+ } finally {
+ segment1.destroy();
+ segment2.destroy();
+ }
+ }
+
+ private IndexSegment createTransformKeySegment(String segmentName, int
colAValue, int colBValue, int numRecords)
+ throws Exception {
+ List<GenericRow> records = new ArrayList<>(numRecords);
+ for (int i = 0; i < numRecords; i++) {
+ GenericRow record = new GenericRow();
+ record.putValue(COL_A, colAValue);
+ record.putValue(COL_B, colBValue);
+ records.add(record);
+ }
+
+ SegmentGeneratorConfig segmentGeneratorConfig =
+ new SegmentGeneratorConfig(TRANSFORM_KEY_TABLE_CONFIG,
TRANSFORM_KEY_SCHEMA);
+ segmentGeneratorConfig.setTableName(RAW_TABLE_NAME);
+ segmentGeneratorConfig.setSegmentName(segmentName);
+ segmentGeneratorConfig.setOutDir(TEMP_DIR.getPath());
+
+ SegmentIndexCreationDriverImpl driver = new
SegmentIndexCreationDriverImpl();
+ driver.init(segmentGeneratorConfig, new GenericRowRecordReader(records));
+ driver.build();
+
+ return ImmutableSegmentLoader.load(new File(TEMP_DIR, segmentName),
ReadMode.mmap);
+ }
+
// ----
// Utils
@BeforeClass
diff --git
a/pinot-core/src/test/java/org/apache/pinot/core/query/utils/OrderByComparatorFactoryTest.java
b/pinot-core/src/test/java/org/apache/pinot/core/query/utils/OrderByComparatorFactoryTest.java
index d0deaf5bbc5..0e88f5a5ee0 100644
---
a/pinot-core/src/test/java/org/apache/pinot/core/query/utils/OrderByComparatorFactoryTest.java
+++
b/pinot-core/src/test/java/org/apache/pinot/core/query/utils/OrderByComparatorFactoryTest.java
@@ -20,13 +20,17 @@
package org.apache.pinot.core.query.utils;
import java.util.Arrays;
+import java.util.Comparator;
import java.util.List;
import java.util.stream.Collectors;
import org.apache.pinot.common.request.context.ExpressionContext;
+import org.apache.pinot.common.request.context.FunctionContext;
import org.apache.pinot.common.request.context.OrderByExpressionContext;
+import org.apache.pinot.core.data.table.Record;
import org.testng.annotations.Test;
import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertTrue;
public class OrderByComparatorFactoryTest {
@@ -104,4 +108,38 @@ public class OrderByComparatorFactoryTest {
assertEquals(extractColumn(_rows, COLUMN2_INDEX), Arrays.asList(1, 2, 3));
}
+
+ private static ExpressionContext function(String name, String arg) {
+ return ExpressionContext.forFunction(
+ new FunctionContext(FunctionContext.Type.TRANSFORM, name,
List.of(ExpressionContext.forIdentifier(arg))));
+ }
+
+ /// getRecordKeyComparator maps each ORDER BY expression to its position in
the GROUP BY list. When several
+ /// group-by keys are transform expressions rather than plain identifiers,
every expression must still resolve to
+ /// its own column, otherwise SortedRecordsMerger sees unequal groups as
equal and merges them.
+ @Test
+ public void testRecordKeyComparatorWithMultipleTransformGroupByKeys() {
+ ExpressionContext key0 = function("datetrunc", "tsColumn");
+ ExpressionContext key1 = function("jsonextractindex", "jsonColumn");
+ List<ExpressionContext> groupByExpressions = List.of(key0, key1);
+ List<OrderByExpressionContext> orderBys =
+ List.of(new OrderByExpressionContext(key0, ASC, NULLS_LAST), new
OrderByExpressionContext(key1, ASC,
+ NULLS_LAST));
+
+ Comparator<Record> comparator =
+ OrderByComparatorFactory.getRecordKeyComparator(orderBys,
groupByExpressions, false);
+
+ // Same second key, different first key: these are distinct groups and
must not compare equal.
+ Record a = new Record(new Object[]{1L, "x", 10.0});
+ Record b = new Record(new Object[]{2L, "x", 20.0});
+ assertTrue(comparator.compare(a, b) < 0, "rows differing only in the first
group key compared equal");
+ assertTrue(comparator.compare(b, a) > 0);
+
+ // Same first key, different second key: also distinct groups.
+ Record c = new Record(new Object[]{1L, "y", 30.0});
+ assertTrue(comparator.compare(a, c) < 0, "rows differing only in the
second group key compared equal");
+
+ // Identical keys are the same group.
+ assertEquals(comparator.compare(a, new Record(new Object[]{1L, "x",
99.0})), 0);
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]