wirybeaver commented on code in PR #19469:
URL: https://github.com/apache/pinot/pull/19469#discussion_r4102820740


##########
pinot-core/src/main/java/org/apache/pinot/core/query/aggregation/function/AnyValueAggregationFunction.java:
##########
@@ -309,7 +312,7 @@ private Object deserializeValue(ByteBuffer buffer) {
       case BIG_DECIMAL:
         return new BigDecimal(new String(deserializeVariableBytes(buffer), 
StandardCharsets.UTF_8));
       case BYTES:
-        return deserializeVariableBytes(buffer);
+        return new ByteArray(deserializeVariableBytes(buffer));

Review Comment:
   The independent, directly tested BYTES fix is in #19667 
(910d0a391295bd82a468c6276cdf5a52b608e373). Its regression test covers raw 
byte[] serialization, ByteArray deserialization, extractFinalResult, and 
reserialization. The original feature commit in #19469 still contains identical 
source lines for now; I have marked #19667 as a prerequisite and will remove 
the duplicate from this PR after the independent fix lands. I am leaving this 
thread open until that cleanup.



##########
pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/operator/AggregateOperator.java:
##########
@@ -268,12 +367,87 @@ private MseBlock.Eos consumeGroupBy() {
     MseBlock block = _input.nextBlock();
     while (block.isData()) {
       _groupByExecutor.processBlock((MseBlock.Data) block);
+      if (_spillThreshold > 0 && _groupByExecutor.getNumGroups() >= 
_spillThreshold) {
+        spillCurrentGroups();
+      }
       checkTerminationAndSampleUsage();
       block = _input.nextBlock();
     }
     return (MseBlock.Eos) block;
   }
 
+  private void spillCurrentGroups() {
+    assert _groupByExecutor != null;
+    if (_groupByExecutor.getNumGroups() == 0) {
+      return;
+    }
+    if (_spillManager == null) {
+      _spillManager =
+          new AggregationSpillManager(_spillPartitions, _groupKeyIds.length, 
getSpillSchema(), _aggFunctions);
+      _spillDirectory = _spillManager.getSpillDirectory();
+    }
+    AggregationSpillManager.SpillResult spillResult =
+        _spillManager.spill(_groupByExecutor.getIntermediateResultIterator());
+    _statMap.merge(StatKey.SPILL_COUNT, 1L);
+    _statMap.merge(StatKey.SPILLED_ROWS, spillResult.getRows());
+    _statMap.merge(StatKey.SPILLED_BYTES, spillResult.getBytes());
+    _groupByExecutor = newInputGroupByExecutor();
+  }
+
+  private DataSchema getSpillSchema() {
+    String[] columnNames = _resultSchema.getColumnNames().clone();
+    ColumnDataType[] columnDataTypes = 
_resultSchema.getColumnDataTypes().clone();
+    int numKeys = _groupKeyIds.length;
+    for (int i = 0; i < _aggFunctions.length; i++) {
+      columnDataTypes[numKeys + i] = _aggFunctions[i].getType() == 
AggregationFunctionType.ANYVALUE

Review Comment:
   Done in 2e40e9402d3069aa86e22cf7e6858b904795b57f. ANY_VALUE now reports 
OBJECT as its intermediate result column type, including when its input type 
has not been resolved on a merge-only executor. AggregateOperator uses the 
function contract directly instead of a function-name special case; spill tests 
cover BYTES, UUID, and other intermediate values.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to