github-actions[bot] commented on code in PR #68500:
URL: https://github.com/apache/doris/pull/68500#discussion_r4226675247


##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -17,27 +17,422 @@
 
 #include "exec/operator/analytic_source_operator.h"
 
+#include <algorithm>
 #include <cstddef>
+#include <cstdint>
+#include <ranges>
 #include <string>
 
 #include "core/column/column_nullable.h"
+#include "core/column/column_vector.h"
 #include "exec/operator/operator.h"
+#include "exec/spill/spill_file.h"
+#include "exec/spill/spill_file_reader.h"
 #include "exprs/vectorized_agg_fn.h"
 
 namespace doris {
 
 AnalyticLocalState::AnalyticLocalState(RuntimeState* state, OperatorXBase* 
parent)
-        : PipelineXLocalState<AnalyticSharedState>(state, parent) {}
+        : Base(state, parent) {}
 
 Status AnalyticLocalState::init(RuntimeState* state, LocalStateInfo& info) {
-    RETURN_IF_ERROR(PipelineXLocalState<AnalyticSharedState>::init(state, 
info));
+    RETURN_IF_ERROR(Base::init(state, info));
     SCOPED_TIMER(exec_time_counter());
     SCOPED_TIMER(_init_timer);
     _get_next_timer = ADD_TIMER(custom_profile(), "GetNextTime");
     _filtered_rows_counter = ADD_COUNTER(custom_profile(), "FilteredRows", 
TUnit::UNIT);
+    _partition_replay_timer = ADD_TIMER(custom_profile(), 
"PartitionReplayTime");
     return Status::OK();
 }
 
+Status AnalyticLocalState::close(RuntimeState* state) {
+    if (_closed) {
+        return Status::OK();
+    }
+    _finish_spill_batch();
+    return Base::close(state);
+}
+
+void AnalyticLocalState::_finish_spill_batch() {
+    if (_batch_reader) {
+        auto st = _batch_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill batch reader 
failed: " << st;
+        _batch_reader.reset();
+    }
+    if (_peer_group_reader) {
+        auto st = _peer_group_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill peer group reader 
failed: " << st;
+        _peer_group_reader.reset();
+    }
+    COUNTER_UPDATE(_memory_used_counter, -_in_memory_batch_bytes);
+    _in_memory_batch_bytes = 0;
+    _current_batch.reset();
+    _replay_block.clear();
+    _replay_block_position = 0;
+    _peer_group_block.clear();
+}
+
+Status AnalyticLocalState::_open_spill_batch(RuntimeState* state,
+                                             
std::shared_ptr<AnalyticSpillBatch> batch) {
+    DCHECK(batch != nullptr);
+    DCHECK_GT(batch->rows, 0);
+    DORIS_CHECK(!batch->partition_ends.empty());
+    DORIS_CHECK_EQ(batch->partition_ends.back(), batch->rows);
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->result_types.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->peer_functions.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->partition_results.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->function_parameters.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->change_to_nullable_flags.size());
+    _current_batch = std::move(batch);
+    _in_memory_block_index = 0;
+    _batch_output_position = 0;
+    _partition_index = 0;
+    _partition_start = 0;
+    _partition_end = _current_batch->partition_ends[0];
+    _peer_group_start = 0;
+    _peer_group_end = 0;
+    _peer_group_block_position = 0;
+    _in_memory_peer_group_index = 0;
+    _peer_group_file_eos = false;
+    for (const auto& block : _current_batch->blocks) {
+        _in_memory_batch_bytes += block.allocated_bytes();
+    }
+    COUNTER_UPDATE(_memory_used_counter, _in_memory_batch_bytes);
+
+    if (_current_batch->data_file) {
+        _batch_reader = _current_batch->data_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_batch_reader->open());
+    }
+    if (_current_batch->peer_group_file) {
+        _peer_group_reader =
+                _current_batch->peer_group_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_peer_group_reader->open());
+    }
+
+    _has_peer_group_function =
+            std::ranges::any_of(_current_batch->strategies, 
[](WindowSpillStrategy strategy) {
+                return strategy == WindowSpillStrategy::PEER_GROUP;
+            });
+    if (_has_peer_group_function) {
+        RETURN_IF_ERROR(_next_peer_group_end(state));
+    }
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_read_batch_block(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    RETURN_IF_CANCELLED(state);
+    if (_batch_reader) {
+        return _batch_reader->read(block, batch_eos);
+    }
+    if (_in_memory_block_index >= _current_batch->blocks.size()) {
+        *batch_eos = true;
+        block->clear();
+        return Status::OK();
+    }
+    auto& next_block = _current_batch->blocks[_in_memory_block_index++];
+    const auto block_bytes = 
static_cast<int64_t>(next_block.allocated_bytes());
+    COUNTER_UPDATE(_memory_used_counter, -block_bytes);
+    _in_memory_batch_bytes -= block_bytes;
+    block->swap(std::move(next_block));
+    *batch_eos = false;
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_next_replay_rows(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    while (_replay_block_position >= _replay_block.rows()) {
+        _replay_block.clear();
+        _replay_block_position = 0;
+        RETURN_IF_ERROR(_read_batch_block(state, &_replay_block, batch_eos));
+        if (*batch_eos) {
+            return Status::OK();
+        }
+    }
+    *batch_eos = false;
+    // Spilled Blocks are coalesced up to the spill buffer size, so a replayed 
Block can be much
+    // larger than the batch size expected by downstream operators.
+    DCHECK_GT(state->batch_size(), 0);
+    const auto batch_size = static_cast<size_t>(state->batch_size());
+    const size_t rows = std::min(batch_size, _replay_block.rows() - 
_replay_block_position);
+    if (_replay_block_position == 0 && rows == _replay_block.rows()) {
+        block->swap(_replay_block);
+        _replay_block.clear();
+        return Status::OK();
+    }
+    Block slice;
+    for (const auto& column : _replay_block) {
+        slice.insert({column.column->cut(_replay_block_position, rows), 
column.type, column.name});
+    }
+    block->swap(slice);
+    _replay_block_position += rows;
+    return Status::OK();
+}
+
+size_t AnalyticLocalState::_spill_replay_reserve_bytes(RuntimeState* state) 
const {
+    if (!_shared_state->spill_enabled.load()) {
+        return 0;
+    }
+    // A spilled record holds Blocks coalesced up to about the spill buffer 
size: opening a reader
+    // allocates the buffer for the largest serialized record of the file and 
every read
+    // deserializes a whole record into a new Block. In-memory batches cost 
nothing to replay.
+    const auto spill_buffer_bytes = 
static_cast<size_t>(state->spill_buffer_size_bytes());
+    if (!_current_batch || _batch_output_position >= _current_batch->rows) {

Review Comment:
   [P2] Reserve for the actual spilled record size. The sink can write a record 
larger than spill_buffer_size_bytes: one wide String row is unsplittable, and 
two ordinary sub-buffer Blocks can coalesce before the post-merge flush. 
SpillFileReader opens a buffer sized to the file's actual maximum record and 
then deserializes it, but this estimates only two configured buffers for a 
data-only batch (and one on later reads). With an 8 MiB buffer, one 
incompressible 32 MiB row already exceeds the 16 MiB opening estimate with the 
read buffer alone. Bound record size at write time or carry its actual maximum 
into replay admission.



##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -17,27 +17,422 @@
 
 #include "exec/operator/analytic_source_operator.h"
 
+#include <algorithm>
 #include <cstddef>
+#include <cstdint>
+#include <ranges>
 #include <string>
 
 #include "core/column/column_nullable.h"
+#include "core/column/column_vector.h"
 #include "exec/operator/operator.h"
+#include "exec/spill/spill_file.h"
+#include "exec/spill/spill_file_reader.h"
 #include "exprs/vectorized_agg_fn.h"
 
 namespace doris {
 
 AnalyticLocalState::AnalyticLocalState(RuntimeState* state, OperatorXBase* 
parent)
-        : PipelineXLocalState<AnalyticSharedState>(state, parent) {}
+        : Base(state, parent) {}
 
 Status AnalyticLocalState::init(RuntimeState* state, LocalStateInfo& info) {
-    RETURN_IF_ERROR(PipelineXLocalState<AnalyticSharedState>::init(state, 
info));
+    RETURN_IF_ERROR(Base::init(state, info));
     SCOPED_TIMER(exec_time_counter());
     SCOPED_TIMER(_init_timer);
     _get_next_timer = ADD_TIMER(custom_profile(), "GetNextTime");
     _filtered_rows_counter = ADD_COUNTER(custom_profile(), "FilteredRows", 
TUnit::UNIT);
+    _partition_replay_timer = ADD_TIMER(custom_profile(), 
"PartitionReplayTime");
     return Status::OK();
 }
 
+Status AnalyticLocalState::close(RuntimeState* state) {
+    if (_closed) {
+        return Status::OK();
+    }
+    _finish_spill_batch();
+    return Base::close(state);
+}
+
+void AnalyticLocalState::_finish_spill_batch() {
+    if (_batch_reader) {
+        auto st = _batch_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill batch reader 
failed: " << st;
+        _batch_reader.reset();
+    }
+    if (_peer_group_reader) {
+        auto st = _peer_group_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill peer group reader 
failed: " << st;
+        _peer_group_reader.reset();
+    }
+    COUNTER_UPDATE(_memory_used_counter, -_in_memory_batch_bytes);
+    _in_memory_batch_bytes = 0;
+    _current_batch.reset();
+    _replay_block.clear();
+    _replay_block_position = 0;
+    _peer_group_block.clear();
+}
+
+Status AnalyticLocalState::_open_spill_batch(RuntimeState* state,
+                                             
std::shared_ptr<AnalyticSpillBatch> batch) {
+    DCHECK(batch != nullptr);
+    DCHECK_GT(batch->rows, 0);
+    DORIS_CHECK(!batch->partition_ends.empty());
+    DORIS_CHECK_EQ(batch->partition_ends.back(), batch->rows);
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->result_types.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->peer_functions.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->partition_results.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->function_parameters.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->change_to_nullable_flags.size());
+    _current_batch = std::move(batch);
+    _in_memory_block_index = 0;
+    _batch_output_position = 0;
+    _partition_index = 0;
+    _partition_start = 0;
+    _partition_end = _current_batch->partition_ends[0];
+    _peer_group_start = 0;
+    _peer_group_end = 0;
+    _peer_group_block_position = 0;
+    _in_memory_peer_group_index = 0;
+    _peer_group_file_eos = false;
+    for (const auto& block : _current_batch->blocks) {
+        _in_memory_batch_bytes += block.allocated_bytes();
+    }
+    COUNTER_UPDATE(_memory_used_counter, _in_memory_batch_bytes);
+
+    if (_current_batch->data_file) {
+        _batch_reader = _current_batch->data_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_batch_reader->open());
+    }
+    if (_current_batch->peer_group_file) {
+        _peer_group_reader =
+                _current_batch->peer_group_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_peer_group_reader->open());
+    }
+
+    _has_peer_group_function =
+            std::ranges::any_of(_current_batch->strategies, 
[](WindowSpillStrategy strategy) {
+                return strategy == WindowSpillStrategy::PEER_GROUP;
+            });
+    if (_has_peer_group_function) {
+        RETURN_IF_ERROR(_next_peer_group_end(state));
+    }
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_read_batch_block(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    RETURN_IF_CANCELLED(state);
+    if (_batch_reader) {
+        return _batch_reader->read(block, batch_eos);
+    }
+    if (_in_memory_block_index >= _current_batch->blocks.size()) {
+        *batch_eos = true;
+        block->clear();
+        return Status::OK();
+    }
+    auto& next_block = _current_batch->blocks[_in_memory_block_index++];
+    const auto block_bytes = 
static_cast<int64_t>(next_block.allocated_bytes());
+    COUNTER_UPDATE(_memory_used_counter, -block_bytes);
+    _in_memory_batch_bytes -= block_bytes;
+    block->swap(std::move(next_block));
+    *batch_eos = false;
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_next_replay_rows(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    while (_replay_block_position >= _replay_block.rows()) {
+        _replay_block.clear();
+        _replay_block_position = 0;
+        RETURN_IF_ERROR(_read_batch_block(state, &_replay_block, batch_eos));
+        if (*batch_eos) {
+            return Status::OK();
+        }
+    }
+    *batch_eos = false;
+    // Spilled Blocks are coalesced up to the spill buffer size, so a replayed 
Block can be much
+    // larger than the batch size expected by downstream operators.
+    DCHECK_GT(state->batch_size(), 0);
+    const auto batch_size = static_cast<size_t>(state->batch_size());
+    const size_t rows = std::min(batch_size, _replay_block.rows() - 
_replay_block_position);
+    if (_replay_block_position == 0 && rows == _replay_block.rows()) {
+        block->swap(_replay_block);
+        _replay_block.clear();
+        return Status::OK();
+    }
+    Block slice;
+    for (const auto& column : _replay_block) {
+        slice.insert({column.column->cut(_replay_block_position, rows), 
column.type, column.name});
+    }
+    block->swap(slice);
+    _replay_block_position += rows;
+    return Status::OK();
+}
+
+size_t AnalyticLocalState::_spill_replay_reserve_bytes(RuntimeState* state) 
const {
+    if (!_shared_state->spill_enabled.load()) {
+        return 0;
+    }
+    // A spilled record holds Blocks coalesced up to about the spill buffer 
size: opening a reader
+    // allocates the buffer for the largest serialized record of the file and 
every read
+    // deserializes a whole record into a new Block. In-memory batches cost 
nothing to replay.
+    const auto spill_buffer_bytes = 
static_cast<size_t>(state->spill_buffer_size_bytes());
+    if (!_current_batch || _batch_output_position >= _current_batch->rows) {
+        // The next call opens the batch at the head of the queue: it opens 
the readers of its
+        // spilled files, reads the first peer group record and then the first 
data record.
+        LockGuard lock(_shared_state->buffer_mutex);
+        if (_shared_state->spill_batches.empty()) {
+            return 0;

Review Comment:
   [P2] Keep the next batch stable across replay admission. After this 
empty-queue check returns zero, the sink can publish a spilled batch before the 
same source call finishes its current batch and dequeues the new one. 
PipelineTask has already completed its memory reservation, so opening that 
batch allocates its reader buffers and first decoded records without admission 
under workload-group pressure. Defer a newly arrived batch to a later 
scheduling round, or reserve from a stable batch snapshot.



##########
be/src/exec/operator/analytic_sink_operator.cpp:
##########
@@ -744,17 +1144,87 @@ Status AnalyticSinkOperatorX::prepare(RuntimeState* 
state) {
                     alignment_of_next_state * alignment_of_next_state;
         }
     }
+
+    _prepare_spill(state);
     return Status::OK();
 }
 
+void AnalyticSinkOperatorX::_prepare_spill(RuntimeState* state) {
+    _window_spill_strategies.resize(_agg_functions_size, 
WindowSpillStrategy::UNSUPPORTED);
+    _window_spill_peer_functions.resize(_agg_functions_size, 
WindowSpillPeerFunction::NONE);
+    // The planner rewrites an OVER clause without ORDER BY to `RANGE BETWEEN 
UNBOUNDED PRECEDING
+    // AND CURRENT ROW`. Every row is then a peer of the current row, so the 
frame still covers
+    // the whole partition and the executor has to buffer the partition until 
it ends.
+    const bool range_to_current_row_without_order =
+            _has_window && _has_range_window && !_has_window_start && 
_has_window_end &&
+            _window.window_end.type == 
TAnalyticWindowBoundaryType::CURRENT_ROW &&
+            _order_by_eq_expr_ctxs.empty();
+    const bool full_partition_frame = !_has_window || (!_has_window_start && 
!_has_window_end) ||
+                                      range_to_current_row_without_order;
+    for (size_t i = 0; i < _agg_functions_size; ++i) {
+        _window_spill_strategies[i] = 
_agg_functions[i]->window_spill_strategy();
+        _window_spill_peer_functions[i] = 
_agg_functions[i]->window_spill_peer_function();
+    }
+    const bool partition_dependent_function =
+            std::ranges::any_of(_window_spill_strategies, 
[](WindowSpillStrategy strategy) {
+                return strategy == WindowSpillStrategy::PARTITION_CARDINALITY 
||
+                       strategy == WindowSpillStrategy::PEER_GROUP;
+            });
+    const bool streaming_frame =
+            !full_partition_frame && _has_window &&
+            (!_has_range_window ||
+             (!_has_window_start && _has_window_end &&
+              _window.window_end.type == 
TAnalyticWindowBoundaryType::CURRENT_ROW));
+    const bool streaming_window = streaming_frame && 
!partition_dependent_function;
+    bool spill_supported = state->enable_spill() && !streaming_window && 
_agg_functions_size > 0;
+    _window_spill_unsupported_reason = "Unsupported";
+    if (!state->enable_spill()) {
+        _window_spill_unsupported_reason = "Disabled";
+    } else if (streaming_window) {
+        _window_spill_unsupported_reason = "Streaming";
+    }
+    for (size_t i = 0; i < _agg_functions_size && spill_supported; ++i) {
+        const auto strategy = _window_spill_strategies[i];
+        bool frame_supported = false;
+        switch (strategy) {
+        case WindowSpillStrategy::PARTITION_REDUCE:
+            frame_supported = full_partition_frame;
+            break;
+        case WindowSpillStrategy::PARTITION_CARDINALITY:
+            frame_supported = _has_window && !_has_range_window && 
!_has_window_start &&
+                              _has_window_end &&
+                              _window.window_end.type == 
TAnalyticWindowBoundaryType::CURRENT_ROW;
+            break;
+        case WindowSpillStrategy::PEER_GROUP:
+            frame_supported = _has_window && _has_range_window && 
!_has_window_start &&
+                              _has_window_end &&
+                              _window.window_end.type == 
TAnalyticWindowBoundaryType::CURRENT_ROW;
+            break;
+        case WindowSpillStrategy::UNSUPPORTED:
+            frame_supported = false;
+            break;
+        }
+        if (!frame_supported) {
+            spill_supported = false;
+            _window_spill_unsupported_reason =
+                    fmt::format("Unsupported: {}", 
_agg_functions[i]->get_name());
+        }
+    }
+    _enable_spill_analytic = spill_supported;
+    _spillable = spill_supported;
+}
+
 Status AnalyticSinkOperatorX::sink_impl(doris::RuntimeState* state, Block* 
input_block, bool eos) {
     auto& local_state = get_local_state(state);
     SCOPED_TIMER(local_state.exec_time_counter());
     COUNTER_UPDATE(local_state.rows_input_counter(), 
(int64_t)input_block->rows());
-    local_state._input_eos = eos;
-    local_state._remove_unused_rows();
     local_state._reserve_mem_size = 0;
     SCOPED_PEAK_MEM(&local_state._reserve_mem_size);
+    if (local_state._spill_enabled) {
+        return local_state._sink_spill(state, input_block, eos);

Review Comment:
   [P2] Reserve spill-write scratch before the first sink call. This new branch 
can force-spill its first ordinary buffer-sized Block, but get_reserve_mem_size 
still returns only the prior call's measured peak, which starts at zero. 
PipelineTask then skips sink memory admission while SpillFileWriter allocates a 
PBlock and serialized/compressed buffers; a peer-sidecar flush can do the same 
without data spill. Under workload-group pressure this can hit the hard limit 
instead of pausing for revocation. Estimate the pending write from the current 
Block and spill settings before entering this branch.



##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -83,6 +488,12 @@ Status 
AnalyticSourceOperatorX::get_block_impl(RuntimeState* state, Block* outpu
     return Status::OK();
 }
 
+size_t AnalyticSourceOperatorX::get_reserve_mem_size(RuntimeState* state) {
+    auto& local_state = get_local_state(state);
+    return OperatorX<AnalyticLocalState>::get_reserve_mem_size(state) +
+           local_state._spill_replay_reserve_bytes(state);

Review Comment:
   [P2] Recompute admission for a retained replay Block. The first read of a 
coalesced spill Block records a large reader-open/deserialization peak in the 
base operator estimate. On later calls that only cut a small batch_size slice 
from _replay_block, this override still requests that prior peak even though 
_spill_replay_reserve_bytes adds zero and no disk read occurs. At a tight 
workload-group watermark the sole consumer can be paused before it can produce 
a smaller estimate. Size the next call from its actual slice and result 
allocations instead of reusing the previous read peak.



##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -17,27 +17,422 @@
 
 #include "exec/operator/analytic_source_operator.h"
 
+#include <algorithm>
 #include <cstddef>
+#include <cstdint>
+#include <ranges>
 #include <string>
 
 #include "core/column/column_nullable.h"
+#include "core/column/column_vector.h"
 #include "exec/operator/operator.h"
+#include "exec/spill/spill_file.h"
+#include "exec/spill/spill_file_reader.h"
 #include "exprs/vectorized_agg_fn.h"
 
 namespace doris {
 
 AnalyticLocalState::AnalyticLocalState(RuntimeState* state, OperatorXBase* 
parent)
-        : PipelineXLocalState<AnalyticSharedState>(state, parent) {}
+        : Base(state, parent) {}
 
 Status AnalyticLocalState::init(RuntimeState* state, LocalStateInfo& info) {
-    RETURN_IF_ERROR(PipelineXLocalState<AnalyticSharedState>::init(state, 
info));
+    RETURN_IF_ERROR(Base::init(state, info));
     SCOPED_TIMER(exec_time_counter());
     SCOPED_TIMER(_init_timer);
     _get_next_timer = ADD_TIMER(custom_profile(), "GetNextTime");
     _filtered_rows_counter = ADD_COUNTER(custom_profile(), "FilteredRows", 
TUnit::UNIT);
+    _partition_replay_timer = ADD_TIMER(custom_profile(), 
"PartitionReplayTime");
     return Status::OK();
 }
 
+Status AnalyticLocalState::close(RuntimeState* state) {
+    if (_closed) {
+        return Status::OK();
+    }
+    _finish_spill_batch();
+    return Base::close(state);
+}
+
+void AnalyticLocalState::_finish_spill_batch() {
+    if (_batch_reader) {
+        auto st = _batch_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill batch reader 
failed: " << st;
+        _batch_reader.reset();
+    }
+    if (_peer_group_reader) {
+        auto st = _peer_group_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill peer group reader 
failed: " << st;
+        _peer_group_reader.reset();
+    }
+    COUNTER_UPDATE(_memory_used_counter, -_in_memory_batch_bytes);
+    _in_memory_batch_bytes = 0;
+    _current_batch.reset();
+    _replay_block.clear();
+    _replay_block_position = 0;
+    _peer_group_block.clear();
+}
+
+Status AnalyticLocalState::_open_spill_batch(RuntimeState* state,
+                                             
std::shared_ptr<AnalyticSpillBatch> batch) {
+    DCHECK(batch != nullptr);
+    DCHECK_GT(batch->rows, 0);
+    DORIS_CHECK(!batch->partition_ends.empty());
+    DORIS_CHECK_EQ(batch->partition_ends.back(), batch->rows);
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->result_types.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->peer_functions.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->partition_results.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->function_parameters.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->change_to_nullable_flags.size());
+    _current_batch = std::move(batch);
+    _in_memory_block_index = 0;
+    _batch_output_position = 0;
+    _partition_index = 0;
+    _partition_start = 0;
+    _partition_end = _current_batch->partition_ends[0];
+    _peer_group_start = 0;
+    _peer_group_end = 0;
+    _peer_group_block_position = 0;
+    _in_memory_peer_group_index = 0;
+    _peer_group_file_eos = false;
+    for (const auto& block : _current_batch->blocks) {
+        _in_memory_batch_bytes += block.allocated_bytes();
+    }
+    COUNTER_UPDATE(_memory_used_counter, _in_memory_batch_bytes);
+
+    if (_current_batch->data_file) {
+        _batch_reader = _current_batch->data_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_batch_reader->open());
+    }
+    if (_current_batch->peer_group_file) {
+        _peer_group_reader =
+                _current_batch->peer_group_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_peer_group_reader->open());
+    }
+
+    _has_peer_group_function =
+            std::ranges::any_of(_current_batch->strategies, 
[](WindowSpillStrategy strategy) {
+                return strategy == WindowSpillStrategy::PEER_GROUP;
+            });
+    if (_has_peer_group_function) {
+        RETURN_IF_ERROR(_next_peer_group_end(state));
+    }
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_read_batch_block(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    RETURN_IF_CANCELLED(state);
+    if (_batch_reader) {
+        return _batch_reader->read(block, batch_eos);
+    }
+    if (_in_memory_block_index >= _current_batch->blocks.size()) {
+        *batch_eos = true;
+        block->clear();
+        return Status::OK();
+    }
+    auto& next_block = _current_batch->blocks[_in_memory_block_index++];
+    const auto block_bytes = 
static_cast<int64_t>(next_block.allocated_bytes());
+    COUNTER_UPDATE(_memory_used_counter, -block_bytes);
+    _in_memory_batch_bytes -= block_bytes;
+    block->swap(std::move(next_block));
+    *batch_eos = false;
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_next_replay_rows(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    while (_replay_block_position >= _replay_block.rows()) {
+        _replay_block.clear();
+        _replay_block_position = 0;
+        RETURN_IF_ERROR(_read_batch_block(state, &_replay_block, batch_eos));
+        if (*batch_eos) {
+            return Status::OK();
+        }
+    }
+    *batch_eos = false;
+    // Spilled Blocks are coalesced up to the spill buffer size, so a replayed 
Block can be much
+    // larger than the batch size expected by downstream operators.
+    DCHECK_GT(state->batch_size(), 0);
+    const auto batch_size = static_cast<size_t>(state->batch_size());
+    const size_t rows = std::min(batch_size, _replay_block.rows() - 
_replay_block_position);
+    if (_replay_block_position == 0 && rows == _replay_block.rows()) {
+        block->swap(_replay_block);
+        _replay_block.clear();
+        return Status::OK();
+    }
+    Block slice;
+    for (const auto& column : _replay_block) {
+        slice.insert({column.column->cut(_replay_block_position, rows), 
column.type, column.name});
+    }
+    block->swap(slice);
+    _replay_block_position += rows;
+    return Status::OK();
+}
+
+size_t AnalyticLocalState::_spill_replay_reserve_bytes(RuntimeState* state) 
const {
+    if (!_shared_state->spill_enabled.load()) {
+        return 0;
+    }
+    // A spilled record holds Blocks coalesced up to about the spill buffer 
size: opening a reader
+    // allocates the buffer for the largest serialized record of the file and 
every read
+    // deserializes a whole record into a new Block. In-memory batches cost 
nothing to replay.
+    const auto spill_buffer_bytes = 
static_cast<size_t>(state->spill_buffer_size_bytes());
+    if (!_current_batch || _batch_output_position >= _current_batch->rows) {
+        // The next call opens the batch at the head of the queue: it opens 
the readers of its
+        // spilled files, reads the first peer group record and then the first 
data record.
+        LockGuard lock(_shared_state->buffer_mutex);
+        if (_shared_state->spill_batches.empty()) {
+            return 0;
+        }
+        const auto& next_batch = _shared_state->spill_batches.front();
+        size_t reserve_bytes = 0;
+        if (next_batch->data_file) {
+            reserve_bytes += 2 * spill_buffer_bytes;
+        }
+        if (next_batch->peer_group_file) {
+            reserve_bytes += 2 * spill_buffer_bytes;
+        }
+        return reserve_bytes;
+    }
+    size_t reserve_bytes = 0;
+    if (_batch_reader && _replay_block_position >= _replay_block.rows()) {
+        reserve_bytes += spill_buffer_bytes;
+    }
+    if (_next_slice_reads_peer_group_record(state)) {
+        reserve_bytes += spill_buffer_bytes;
+    }
+    return reserve_bytes;
+}
+
+bool AnalyticLocalState::_next_slice_reads_peer_group_record(RuntimeState* 
state) const {
+    if (!_peer_group_reader) {
+        return false;
+    }
+    // The next call outputs at most batch_size rows starting at 
_batch_output_position and needs
+    // a new peer group record once one of them is not covered by the last 
group end read so far.
+    DCHECK_GT(_peer_group_block.rows(), 0);
+    const auto& column =
+            assert_cast<const 
ColumnInt64&>(*_peer_group_block.get_by_position(0).column);
+    const int64_t last_peer_group_end = column.get_data().back();
+    const int64_t slice_end =
+            std::min<int64_t>(_batch_output_position + state->batch_size(), 
_current_batch->rows);
+    return last_peer_group_end < slice_end;

Review Comment:
   [P2] Limit sidecar lookahead to the next data Block. This uses batch_size 
even when _next_replay_rows will return fewer rows at the current data Block 
boundary. If the loaded sidecar record ends at that boundary, the next call 
needs no sidecar read, yet it reserves one full spill buffer and can pause the 
only consumer at the workload-group high watermark. Use the actual next output 
row count when deciding whether another peer record will be read.



##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -17,27 +17,422 @@
 
 #include "exec/operator/analytic_source_operator.h"
 
+#include <algorithm>
 #include <cstddef>
+#include <cstdint>
+#include <ranges>
 #include <string>
 
 #include "core/column/column_nullable.h"
+#include "core/column/column_vector.h"
 #include "exec/operator/operator.h"
+#include "exec/spill/spill_file.h"
+#include "exec/spill/spill_file_reader.h"
 #include "exprs/vectorized_agg_fn.h"
 
 namespace doris {
 
 AnalyticLocalState::AnalyticLocalState(RuntimeState* state, OperatorXBase* 
parent)
-        : PipelineXLocalState<AnalyticSharedState>(state, parent) {}
+        : Base(state, parent) {}
 
 Status AnalyticLocalState::init(RuntimeState* state, LocalStateInfo& info) {
-    RETURN_IF_ERROR(PipelineXLocalState<AnalyticSharedState>::init(state, 
info));
+    RETURN_IF_ERROR(Base::init(state, info));
     SCOPED_TIMER(exec_time_counter());
     SCOPED_TIMER(_init_timer);
     _get_next_timer = ADD_TIMER(custom_profile(), "GetNextTime");
     _filtered_rows_counter = ADD_COUNTER(custom_profile(), "FilteredRows", 
TUnit::UNIT);
+    _partition_replay_timer = ADD_TIMER(custom_profile(), 
"PartitionReplayTime");
     return Status::OK();
 }
 
+Status AnalyticLocalState::close(RuntimeState* state) {
+    if (_closed) {
+        return Status::OK();
+    }
+    _finish_spill_batch();
+    return Base::close(state);
+}
+
+void AnalyticLocalState::_finish_spill_batch() {
+    if (_batch_reader) {
+        auto st = _batch_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill batch reader 
failed: " << st;
+        _batch_reader.reset();
+    }
+    if (_peer_group_reader) {
+        auto st = _peer_group_reader->close();
+        LOG_IF(WARNING, !st.ok()) << "close analytic spill peer group reader 
failed: " << st;
+        _peer_group_reader.reset();
+    }
+    COUNTER_UPDATE(_memory_used_counter, -_in_memory_batch_bytes);
+    _in_memory_batch_bytes = 0;
+    _current_batch.reset();
+    _replay_block.clear();
+    _replay_block_position = 0;
+    _peer_group_block.clear();
+}
+
+Status AnalyticLocalState::_open_spill_batch(RuntimeState* state,
+                                             
std::shared_ptr<AnalyticSpillBatch> batch) {
+    DCHECK(batch != nullptr);
+    DCHECK_GT(batch->rows, 0);
+    DORIS_CHECK(!batch->partition_ends.empty());
+    DORIS_CHECK_EQ(batch->partition_ends.back(), batch->rows);
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->result_types.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->peer_functions.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), batch->partition_results.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->function_parameters.size());
+    DORIS_CHECK_EQ(batch->strategies.size(), 
batch->change_to_nullable_flags.size());
+    _current_batch = std::move(batch);
+    _in_memory_block_index = 0;
+    _batch_output_position = 0;
+    _partition_index = 0;
+    _partition_start = 0;
+    _partition_end = _current_batch->partition_ends[0];
+    _peer_group_start = 0;
+    _peer_group_end = 0;
+    _peer_group_block_position = 0;
+    _in_memory_peer_group_index = 0;
+    _peer_group_file_eos = false;
+    for (const auto& block : _current_batch->blocks) {
+        _in_memory_batch_bytes += block.allocated_bytes();
+    }
+    COUNTER_UPDATE(_memory_used_counter, _in_memory_batch_bytes);
+
+    if (_current_batch->data_file) {
+        _batch_reader = _current_batch->data_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_batch_reader->open());
+    }
+    if (_current_batch->peer_group_file) {
+        _peer_group_reader =
+                _current_batch->peer_group_file->create_reader(state, 
operator_profile());
+        RETURN_IF_ERROR(_peer_group_reader->open());
+    }
+
+    _has_peer_group_function =
+            std::ranges::any_of(_current_batch->strategies, 
[](WindowSpillStrategy strategy) {
+                return strategy == WindowSpillStrategy::PEER_GROUP;
+            });
+    if (_has_peer_group_function) {
+        RETURN_IF_ERROR(_next_peer_group_end(state));
+    }
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_read_batch_block(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    RETURN_IF_CANCELLED(state);
+    if (_batch_reader) {
+        return _batch_reader->read(block, batch_eos);
+    }
+    if (_in_memory_block_index >= _current_batch->blocks.size()) {
+        *batch_eos = true;
+        block->clear();
+        return Status::OK();
+    }
+    auto& next_block = _current_batch->blocks[_in_memory_block_index++];
+    const auto block_bytes = 
static_cast<int64_t>(next_block.allocated_bytes());
+    COUNTER_UPDATE(_memory_used_counter, -block_bytes);
+    _in_memory_batch_bytes -= block_bytes;
+    block->swap(std::move(next_block));
+    *batch_eos = false;
+    return Status::OK();
+}
+
+Status AnalyticLocalState::_next_replay_rows(RuntimeState* state, Block* 
block, bool* batch_eos) {
+    while (_replay_block_position >= _replay_block.rows()) {
+        _replay_block.clear();
+        _replay_block_position = 0;
+        RETURN_IF_ERROR(_read_batch_block(state, &_replay_block, batch_eos));
+        if (*batch_eos) {
+            return Status::OK();
+        }
+    }
+    *batch_eos = false;
+    // Spilled Blocks are coalesced up to the spill buffer size, so a replayed 
Block can be much
+    // larger than the batch size expected by downstream operators.
+    DCHECK_GT(state->batch_size(), 0);
+    const auto batch_size = static_cast<size_t>(state->batch_size());
+    const size_t rows = std::min(batch_size, _replay_block.rows() - 
_replay_block_position);
+    if (_replay_block_position == 0 && rows == _replay_block.rows()) {
+        block->swap(_replay_block);
+        _replay_block.clear();
+        return Status::OK();
+    }
+    Block slice;
+    for (const auto& column : _replay_block) {
+        slice.insert({column.column->cut(_replay_block_position, rows), 
column.type, column.name});
+    }
+    block->swap(slice);
+    _replay_block_position += rows;
+    return Status::OK();
+}
+
+size_t AnalyticLocalState::_spill_replay_reserve_bytes(RuntimeState* state) 
const {
+    if (!_shared_state->spill_enabled.load()) {
+        return 0;
+    }
+    // A spilled record holds Blocks coalesced up to about the spill buffer 
size: opening a reader
+    // allocates the buffer for the largest serialized record of the file and 
every read
+    // deserializes a whole record into a new Block. In-memory batches cost 
nothing to replay.
+    const auto spill_buffer_bytes = 
static_cast<size_t>(state->spill_buffer_size_bytes());
+    if (!_current_batch || _batch_output_position >= _current_batch->rows) {
+        // The next call opens the batch at the head of the queue: it opens 
the readers of its
+        // spilled files, reads the first peer group record and then the first 
data record.
+        LockGuard lock(_shared_state->buffer_mutex);
+        if (_shared_state->spill_batches.empty()) {
+            return 0;
+        }
+        const auto& next_batch = _shared_state->spill_batches.front();
+        size_t reserve_bytes = 0;
+        if (next_batch->data_file) {
+            reserve_bytes += 2 * spill_buffer_bytes;
+        }
+        if (next_batch->peer_group_file) {
+            reserve_bytes += 2 * spill_buffer_bytes;
+        }
+        return reserve_bytes;
+    }
+    size_t reserve_bytes = 0;
+    if (_batch_reader && _replay_block_position >= _replay_block.rows()) {
+        reserve_bytes += spill_buffer_bytes;
+    }
+    if (_next_slice_reads_peer_group_record(state)) {
+        reserve_bytes += spill_buffer_bytes;
+    }
+    return reserve_bytes;
+}
+
+bool AnalyticLocalState::_next_slice_reads_peer_group_record(RuntimeState* 
state) const {
+    if (!_peer_group_reader) {
+        return false;
+    }
+    // The next call outputs at most batch_size rows starting at 
_batch_output_position and needs
+    // a new peer group record once one of them is not covered by the last 
group end read so far.
+    DCHECK_GT(_peer_group_block.rows(), 0);
+    const auto& column =
+            assert_cast<const 
ColumnInt64&>(*_peer_group_block.get_by_position(0).column);
+    const int64_t last_peer_group_end = column.get_data().back();
+    const int64_t slice_end =
+            std::min<int64_t>(_batch_output_position + state->batch_size(), 
_current_batch->rows);
+    return last_peer_group_end < slice_end;
+}
+
+Status AnalyticLocalState::_next_peer_group_end(RuntimeState* state) {
+    RETURN_IF_CANCELLED(state);
+    _peer_group_start = _peer_group_end;
+    if (!_peer_group_reader) {
+        DORIS_CHECK_LT(_in_memory_peer_group_index, 
_current_batch->peer_group_ends.size());
+        _peer_group_end = 
_current_batch->peer_group_ends[_in_memory_peer_group_index++];
+    } else {
+        while (_peer_group_block_position >= _peer_group_block.rows()) {
+            _peer_group_block.clear();
+            RETURN_IF_ERROR(_peer_group_reader->read(&_peer_group_block, 
&_peer_group_file_eos));
+            _peer_group_block_position = 0;
+            DORIS_CHECK(!_peer_group_file_eos || !_peer_group_block.empty());
+        }
+        const auto& column =
+                assert_cast<const 
ColumnInt64&>(*_peer_group_block.get_by_position(0).column);
+        _peer_group_end = column.get_data()[_peer_group_block_position++];
+    }
+    DORIS_CHECK_GT(_peer_group_end, _peer_group_start);
+    DORIS_CHECK_LE(_peer_group_end, _current_batch->rows);
+    return Status::OK();
+}
+
+void AnalyticLocalState::_next_spill_partition() {
+    ++_partition_index;
+    DORIS_CHECK_LT(_partition_index, _current_batch->partition_ends.size());
+    _partition_start = _partition_end;
+    _partition_end = _current_batch->partition_ends[_partition_index];
+    DORIS_CHECK_GT(_partition_end, _partition_start);
+}
+
+Status AnalyticLocalState::_append_spill_results(RuntimeState* state, Block* 
block) {
+    const auto rows = block->rows();
+    DCHECK_GT(rows, 0);
+    DCHECK_LE(_batch_output_position + static_cast<int64_t>(rows), 
_current_batch->rows);
+    const auto function_count = _current_batch->strategies.size();
+    std::vector<MutableColumnPtr> results(function_count);
+    std::vector<IColumn*> computed_result_columns(function_count, nullptr);
+    _create_spill_result_columns(rows, results, computed_result_columns);
+
+    size_t offset = 0;
+    while (offset < rows) {
+        while (_batch_output_position >= _partition_end) {
+            _next_spill_partition();
+        }
+        const size_t segment_rows =
+                std::min<size_t>(rows - offset, _partition_end - 
_batch_output_position);
+        _append_partition_results(segment_rows, computed_result_columns, 
results);
+        if (_has_peer_group_function) {
+            RETURN_IF_ERROR(
+                    _append_peer_group_results(state, segment_rows, 
computed_result_columns));
+        }
+        _batch_output_position += segment_rows;
+        offset += segment_rows;
+    }
+    _insert_spill_result_columns(block, results);
+    return Status::OK();
+}
+
+void AnalyticLocalState::_create_spill_result_columns(
+        size_t rows, std::vector<MutableColumnPtr>& results,
+        std::vector<IColumn*>& computed_result_columns) const {
+    const auto function_count = _current_batch->strategies.size();
+    for (size_t i = 0; i < function_count; ++i) {
+        results[i] = _current_batch->result_types[i]->create_column();
+        results[i]->reserve(rows);
+        if (_current_batch->strategies[i] == 
WindowSpillStrategy::PARTITION_CARDINALITY ||

Review Comment:
   [P2] Admit the first in-memory replay for its result columns. A queued batch 
with no spill files adds zero to source reservation, and the first call gets 
only the 4 MiB base minimum. With the legal batch_size of 65,535, one narrow 
input Block can remain in memory while 32 Int64 window results reserve about 16 
MiB here before insertion. Under a tight memory limit this bypasses admission 
and can hit the hard limit. Include next-slice row count and result 
widths/function count in the source estimate even for in-memory batches.



##########
be/src/exec/operator/analytic_sink_operator.cpp:
##########
@@ -908,6 +1378,25 @@ size_t 
AnalyticSinkOperatorX::get_reserve_mem_size(RuntimeState* state, bool eos
     return local_state._reserve_mem_size;
 }
 
+size_t AnalyticSinkOperatorX::revocable_mem_size(RuntimeState* state) const {
+    const auto& local_state = get_local_state(state);
+    if (!local_state._spill_enabled || !local_state._batch_store) {
+        return 0;
+    }
+    const auto bytes = local_state._batch_store->revocable_mem_size();
+    const auto min_revocable_mem = 
static_cast<size_t>(state->spill_min_revocable_mem());
+    return bytes > min_revocable_mem ? bytes : 0;
+}

Review Comment:
   [P2] Expose staged spill bytes to revocation below 32 MiB. After a partition 
begins spilling, the write buffer and peer-end buffer can retain several MiB 
yet stay below the default spill_min_revocable_mem threshold, so this returns 
zero. PipelineTask would revoke at 512 KiB, and spill() can flush those 
buffers. If workload-group admission then pauses this sink while the partition 
remains open, no batch is sealed for the source and the pause handler finds no 
revocable task; it can time out. Report flushable bytes at least down to the 
pipeline's revoke floor under pressure.



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