github-actions[bot] commented on code in PR #68500:
URL: https://github.com/apache/doris/pull/68500#discussion_r4139868859
##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -17,27 +17,387 @@
#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() || _replay_block_position <
_replay_block.rows()) {
+ return 0;
+ }
+ // The next call reads a new Block. A spilled record holds Blocks
coalesced up to about the
+ // spill buffer size and is deserialized into a new Block. Opening the
reader of the next
+ // batch additionally allocates the buffer for its largest serialized
record, and it is not
+ // known yet whether that batch was spilled.
+ const auto spill_buffer_bytes =
static_cast<size_t>(state->spill_buffer_size_bytes());
+ if (_current_batch && _batch_output_position < _current_batch->rows) {
+ return _current_batch->data_file ? spill_buffer_bytes : 0;
+ }
+ return 2 * spill_buffer_bytes;
Review Comment:
[P2] Inspect the queued batch before reserving spill-reader memory. This
returns two spill buffers at every batch boundary, even when the next batch was
kept entirely in memory below the sink spill threshold. After publication its
bytes are no longer sink-revocable, yet PipelineTask must reserve these unused
reader buffers before `_get_spill_block()` can dequeue it. If a workload group
rejects that reservation at its high watermark while the query has no revocable
task, the only consumer can remain paused and hit the paused-queue timeout
although replay would open no file. Base the reservation on the queued batch's
`data_file` state, or otherwise let in-memory batches be consumed without a
disk-read reservation.
##########
be/src/exec/operator/analytic_source_operator.cpp:
##########
@@ -17,27 +17,387 @@
#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) {
Review Comment:
[P2] Include the peer-group sidecar in replay memory reservation. A batch
can have both `data_file` and `peer_group_file`; this opens both readers and
`_next_peer_group_end()` retains a decoded sidecar Block before the same
`get_block()` reads the data Block. Each reader also retains its own serialized
buffer/PBlock, while `_spill_replay_reserve_bytes()` adds only two spill
buffers for opening any batch. With many distinct peers, the sidecar Block
alone can approach `spill_buffer_size_bytes()`, so the first replay can
allocate beyond its reservation and hit normal memory limits. Account for the
sidecar reader and Block when reserving replay memory.
--
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]