HappenLee commented on code in PR #68500:
URL: https://github.com/apache/doris/pull/68500#discussion_r4120354529


##########
be/src/exec/operator/analytic_spill.cpp:
##########
@@ -0,0 +1,154 @@
+// 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.
+
+#include "exec/operator/analytic_spill.h"
+
+#include <fmt/format.h>
+
+#include <utility>
+
+#include "core/column/column_vector.h"
+#include "core/data_type/data_type_number.h"
+#include "exec/spill/spill_file.h"
+#include "exec/spill/spill_file_manager.h"
+#include "exec/spill/spill_file_writer.h"
+#include "runtime/exec_env.h"
+#include "runtime/runtime_profile.h"
+#include "runtime/runtime_state.h"
+
+namespace doris {
+
+AnalyticSpillBatchStore::AnalyticSpillBatchStore(RuntimeProfile* profile, int 
node_id,
+                                                 bool has_peer_groups)
+        : _profile(profile), _node_id(node_id), 
_has_peer_groups(has_peer_groups) {}
+
+AnalyticSpillBatchStore::~AnalyticSpillBatchStore() {
+    if (_peer_group_writer) {
+        (void)_peer_group_writer->close();
+    }
+    if (_data_writer) {
+        (void)_data_writer->close();
+    }
+}
+
+Status AnalyticSpillBatchStore::_create_writer(RuntimeState* state, const 
char* label,
+                                               SpillFileSPtr& file, 
SpillFileWriterSPtr& writer) {
+    auto relative_path =
+            fmt::format("{}/{}-{}-{}-{}", print_id(state->query_id()), label, 
_node_id,
+                        state->task_id(), 
ExecEnv::GetInstance()->spill_file_mgr()->next_id());
+    RETURN_IF_ERROR(
+            
ExecEnv::GetInstance()->spill_file_mgr()->create_spill_file(relative_path, 
file));
+    return file->create_writer(state, _profile, writer);
+}
+
+Status AnalyticSpillBatchStore::append_block(RuntimeState* state, Block block) 
{
+    DCHECK_GT(block.rows(), 0);
+    _rows += block.rows();
+    if (_data_writer) {
+        return _data_writer->write_block(state, block);

Review Comment:
   **[P2] Coalesce small blocks before writing analytic spill, following the 
hash join spill pattern**
   
   Once `_data_writer` exists, every incoming Block is immediately passed to 
`write_block()`, which serializes and ZSTD-compresses it and appends it to the 
file. `_append_spill_rows()` only splits oversized input blocks; it does not 
combine smaller ones. The initial `spill()` loop also writes each retained 
Block separately.
   
   For a large window partition arriving in 128 KiB blocks with an 8 MiB spill 
buffer, this produces roughly 64 serialization/compression/append operations 
per 8 MiB of input, instead of one coalesced write. This is a call-count 
comparison, not a measured 64x slowdown; the extra per-block work can limit 
spill throughput for narrow rows or small input batches.
   
   We already have a suitable pattern in the hash join implementation:
   
   - 
[`_partition_block()`](https://github.com/apache/doris/blob/83fdbf977538e19ffcce1bda6d0ea3b1ebdb360e/be/src/exec/operator/partitioned_hash_join_sink_operator.cpp#L342-L373)
 accumulates rows in a `MutableBlock`.
   - 
[`_execute_spill_partitioned_blocks()`](https://github.com/apache/doris/blob/83fdbf977538e19ffcce1bda6d0ea3b1ebdb360e/be/src/exec/operator/partitioned_hash_join_sink_operator.cpp#L297-L318)
 flushes when `allocated_bytes() >= spill_buffer_size_bytes()` or a flush is 
forced.
   - 
[`revoke_memory()`](https://github.com/apache/doris/blob/83fdbf977538e19ffcce1bda6d0ea3b1ebdb360e/be/src/exec/operator/partitioned_hash_join_sink_operator.cpp#L323-L338)
 forces the remaining buffers out so reclamation can make progress.
   
   Please consider the same bounded buffering pattern here: combine small 
payload blocks in input order, flush at the size threshold, and flush the tail 
on revoke and before sealing/closing the writer. Include the buffer in memory 
tracking and `revocable_mem_size()`; already large blocks can retain the 
direct-write path. This can stay within the batch store without changing 
aggregate-update boundaries or peer offsets.
   
   A focused test should cover multiple small appends, the final partial 
buffer, and forced reclamation. A small-block benchmark comparing 
`SpillWriteBlockCount`, serialization time, throughput, and peak memory would 
quantify the tradeoff.



-- 
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