wirybeaver commented on code in PR #19469:
URL: https://github.com/apache/pinot/pull/19469#discussion_r4102822281
##########
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) {
Review Comment:
Done in 2e40e9402d3069aa86e22cf7e6858b904795b57f. The policy is now
MultistageGroupByExecutor.shouldSpill(), invoked by AggregateOperator after an
input block. A future byte-based trigger can refine both the executor-owned
predicate and its evaluation frequency without embedding new policy in
AggregateOperator; that follow-up is #19666.
[addressed by agent]
##########
pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java:
##########
@@ -808,6 +808,20 @@ public static class QueryOptionKey {
/// Flush threshold for streaming group-by on MSE leaf stages.
public static final String STREAMING_GROUP_BY_FLUSH_THRESHOLD =
"streamingGroupByFlushThreshold";
+ /// Maximum number of groups retained by a keyed MSE aggregation
executor before its intermediate states are
+ /// spilled to local disk. This option is honored only when the
server-level aggregation spill gate is enabled.
+ /// An absent value disables spilling. The first spill version does
not apply to global aggregation,
+ /// leaf-final-result, or group-trim modes.
+ public static final String MSE_AGGREGATION_SPILL_THRESHOLD =
"mseAggregationSpillThreshold";
Review Comment:
Done in 2e40e9402d3069aa86e22cf7e6858b904795b57f. The option javadoc
explicitly calls group count a v1 proxy for retained memory and links the
byte-trigger follow-up #19666. That issue tracks per-function size estimation,
within-block evaluation, and interaction with the existing group limit.
[addressed by agent]
##########
pinot-common/src/main/java/org/apache/pinot/common/utils/config/QueryOptionsUtils.java:
##########
@@ -549,6 +549,29 @@ public static Integer
getStreamingGroupByFlushThreshold(Map<String, String> quer
return
checkedParseIntNonNegative(QueryOptionKey.STREAMING_GROUP_BY_FLUSH_THRESHOLD,
value);
}
+ @Nullable
+ public static Integer getMSEAggregationSpillThreshold(Map<String, String>
queryOptions) {
Review Comment:
Done in 2e40e9402d3069aa86e22cf7e6858b904795b57f. The unreleased query
option is renamed to mseAggregationSpillMaxGroups, leaving room for a future
mseAggregationSpillMaxBytes option (#19666). numGroupsLimit continues to cap
each input table even when maxGroups is set higher.
[addressed by agent]
##########
pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/operator/AggregationSpillManager.java:
##########
@@ -0,0 +1,428 @@
+/**
+ * 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.query.runtime.operator;
+
+import com.google.common.annotations.VisibleForTesting;
+import java.io.BufferedInputStream;
+import java.io.DataInputStream;
+import java.io.IOException;
+import java.io.UncheckedIOException;
+import java.nio.ByteBuffer;
+import java.nio.channels.FileChannel;
+import java.nio.file.FileVisitResult;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.nio.file.SimpleFileVisitor;
+import java.nio.file.StandardOpenOption;
+import java.nio.file.attribute.BasicFileAttributes;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.function.Consumer;
+import org.apache.pinot.common.datablock.DataBlock;
+import org.apache.pinot.common.datablock.DataBlockUtils;
+import org.apache.pinot.common.utils.DataSchema;
+import org.apache.pinot.core.query.aggregation.function.AggregationFunction;
+import org.apache.pinot.query.runtime.blocks.MseBlock;
+import org.apache.pinot.query.runtime.blocks.RowHeapDataBlock;
+import org.apache.pinot.query.runtime.blocks.SerializedDataBlock;
+import org.apache.pinot.spi.query.QueryThreadContext;
+import org.apache.pinot.spi.utils.CommonConstants.Server;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/// Manages hash-partitioned aggregation spill files for one operator. The
caller must finish reading before calling
+/// [#close()], which recursively removes the operator-scoped directory. This
class is not thread-safe.
+@SuppressWarnings("rawtypes")
+class AggregationSpillManager implements AutoCloseable {
Review Comment:
Done in the #19469 description alongside
2e40e9402d3069aa86e22cf7e6858b904795b57f. Changed Close: #12080 to Partially
addresses #12080 and explicitly left SSE spill and parallel combine open. This
PR remains MSE-scoped; sharing its single-threaded spill manager with SSE is
not included because the SSE combine stage has different concurrency and
attachment points.
[addressed by agent]
--
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]