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]