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]

Reply via email to