github-actions[bot] commented on code in PR #68500:
URL: https://github.com/apache/doris/pull/68500#discussion_r4226970903
##########
be/src/exec/operator/analytic_sink_operator.cpp:
##########
@@ -631,6 +668,369 @@ int64_t
AnalyticSinkLocalState::find_first_not_equal(IColumn* reference_column,
return end - 1;
}
+bool AnalyticSinkLocalState::_keys_equal(const std::vector<ColumnPtr>& lhs,
size_t lhs_row,
+ const std::vector<ColumnPtr>& rhs,
size_t rhs_row) const {
+ DCHECK_EQ(lhs.size(), rhs.size());
+ for (size_t i = 0; i < lhs.size(); ++i) {
+ if (lhs[i]->compare_at(lhs_row, rhs_row, *rhs[i], 1) != 0) {
+ return false;
+ }
+ }
+ return true;
+}
+
+void AnalyticSinkLocalState::_save_last_keys(const std::vector<ColumnPtr>&
columns, size_t row,
+ std::vector<ColumnPtr>&
last_keys) {
+ last_keys.clear();
+ last_keys.reserve(columns.size());
+ for (const auto& column : columns) {
+ auto last_value = column->clone_empty();
+ last_value->insert_from(*column, row);
+ last_keys.emplace_back(std::move(last_value));
+ }
+}
+
+Status AnalyticSinkLocalState::_materialize_spill_columns(
+ Block* input_block, std::vector<std::vector<ColumnPtr>>* agg_columns,
+ std::vector<ColumnPtr>* partition_columns, std::vector<ColumnPtr>*
order_columns) {
+ const auto original_columns = input_block->columns();
+ agg_columns->resize(_agg_functions_size);
+ for (size_t i = 0; i < _agg_functions_size; ++i) {
+ (*agg_columns)[i].reserve(_agg_expr_ctxs[i].size());
+ for (const auto& expr : _agg_expr_ctxs[i]) {
+ ColumnPtr column;
+ RETURN_IF_ERROR(expr->execute(input_block, column));
+
(*agg_columns)[i].emplace_back(column->convert_to_full_column_if_const());
+ }
+ }
+ partition_columns->reserve(_partition_by_eq_expr_ctxs.size());
+ for (const auto& expr : _partition_by_eq_expr_ctxs) {
+ ColumnPtr column;
+ RETURN_IF_ERROR(expr->execute(input_block, column));
+
partition_columns->emplace_back(column->convert_to_full_column_if_const());
+ }
+ order_columns->reserve(_order_by_eq_expr_ctxs.size());
+ for (const auto& expr : _order_by_eq_expr_ctxs) {
+ ColumnPtr column;
+ RETURN_IF_ERROR(expr->execute(input_block, column));
+ order_columns->emplace_back(column->convert_to_full_column_if_const());
+ }
+ Block::erase_useless_column(input_block, original_columns);
+ return Status::OK();
+}
+
+void AnalyticSinkLocalState::_update_spill_aggregate_states(
+ size_t start, size_t length, const
std::vector<std::vector<ColumnPtr>>& agg_columns) {
+ for (size_t i = 0; i < _agg_functions_size; ++i) {
+ if (_spill_strategies[i] != WindowSpillStrategy::PARTITION_REDUCE) {
+ continue;
+ }
+ std::vector<const IColumn*> columns;
+ columns.reserve(agg_columns[i].size());
+ for (const auto& column : agg_columns[i]) {
+ columns.push_back(column.get());
+ }
+ _agg_functions[i]->add_range_single_place(0, start + length, start,
start + length,
+ _fn_place_ptr +
_offsets_of_aggregate_states[i],
+ columns.data(),
_shared_state->agg_arena_pool,
+ &_use_null_result[i],
+
&_could_use_previous_result[i]);
+ }
+}
+
+Status AnalyticSinkLocalState::_record_peer_groups(RuntimeState* state,
+ const
std::vector<ColumnPtr>& order_columns,
+ size_t start, size_t
length, int64_t batch_row) {
+ if (!_has_peer_group_functions || order_columns.empty()) {
+ return Status::OK();
+ }
+ for (size_t offset = 0; offset < length; ++offset) {
+ if ((offset & 4095) == 0) {
+ RETURN_IF_CANCELLED(state);
+ }
+ const size_t row = start + offset;
+ const bool same_group = offset == 0
+ ? (_open_partition_rows == 0 ||
+ _keys_equal(_last_order_keys, 0,
order_columns, row))
+ : _keys_equal(order_columns, row - 1,
order_columns, row);
+ if (!same_group) {
+ RETURN_IF_ERROR(_batch_store->append_peer_group_end(state,
batch_row + offset));
+ }
+ }
+ return Status::OK();
+}
+
+Status AnalyticSinkLocalState::_accumulate_spill_range(
+ RuntimeState* state, const std::vector<std::vector<ColumnPtr>>&
agg_columns,
+ const std::vector<ColumnPtr>& order_columns, size_t start, size_t
length,
+ int64_t batch_row) {
+ DCHECK_GT(length, 0);
+ DCHECK_LE(batch_row + static_cast<int64_t>(length), _batch_store->rows());
+ if (_open_partition_rows == 0) {
+ for (size_t i = 0; i < _agg_functions_size; ++i) {
+ if (_spill_strategies[i] ==
WindowSpillStrategy::PARTITION_CARDINALITY) {
+ DORIS_CHECK_EQ(agg_columns[i].size(), 1);
+ _open_partition_parameters[i] =
agg_columns[i][0]->get_int(start);
+ DORIS_CHECK_GT(_open_partition_parameters[i], 0);
+ }
+ }
+ }
+ _update_spill_aggregate_states(start, length, agg_columns);
+ RETURN_IF_ERROR(_record_peer_groups(state, order_columns, start, length,
batch_row));
+ _open_partition_rows += length;
+ _open_partition_end = batch_row + length;
+ return Status::OK();
+}
+
+Status AnalyticSinkLocalState::_append_spill_rows(RuntimeState* state, Block*
input_block,
+ size_t start, size_t length)
{
+ DCHECK_GT(length, 0);
+ const size_t input_rows = input_block->rows();
+ const size_t average_row_bytes = std::max<size_t>(1, input_block->bytes()
/ input_rows);
+ const size_t rows_per_block = std::max<size_t>(
+ 1, static_cast<size_t>(state->spill_buffer_size_bytes()) /
average_row_bytes);
+ size_t offset = 0;
+ while (offset < length) {
+ RETURN_IF_CANCELLED(state);
+ const size_t rows = std::min(rows_per_block, length - offset);
+ if (rows == input_rows) {
+ // The whole input Block belongs to this range and needs no split:
take it over.
+ RETURN_IF_ERROR(_batch_store->append_block(state,
std::move(*input_block)));
+ } else {
+ MutableBlock mutable_block(input_block->clone_empty());
+ RETURN_IF_ERROR(mutable_block.add_rows(input_block, start +
offset, rows));
+ RETURN_IF_ERROR(_batch_store->append_block(state,
mutable_block.to_block()));
+ }
+ _update_spill_memory_usage();
+ // Once spilled, the store bounds its own write buffer by the spill
buffer size. Forcing a
+ // spill for every append would flush that buffer and write every
small Block separately.
+ if (!_batch_store->is_spilled() &&
+ (state->enable_force_spill() ||
+ std::cmp_greater_equal(_batch_store->revocable_mem_size(),
+
state->spill_analytic_sink_mem_limit_bytes()))) {
+ RETURN_IF_ERROR(_spill_batch_store(state));
+ }
+ offset += rows;
+ }
+ return Status::OK();
+}
+
+Status AnalyticSinkLocalState::_spill_batch_store(RuntimeState* state) {
+ RETURN_IF_ERROR(_batch_store->spill(state));
+ _update_spill_memory_usage();
+ custom_profile()->add_info_string("WindowSpillMode", "Spilled");
+ return Status::OK();
+}
+
+void AnalyticSinkLocalState::_update_spill_memory_usage() {
+ const size_t current_bytes = _batch_store ?
_batch_store->revocable_mem_size() : 0;
+ const auto delta =
+ static_cast<int64_t>(current_bytes) -
static_cast<int64_t>(_batch_buffered_bytes);
+ _peak_partition_buffered_bytes->add(delta);
+ COUNTER_UPDATE(_memory_used_counter, delta);
+ COUNTER_UPDATE(_blocks_memory_usage, delta);
+ _batch_buffered_bytes = current_bytes;
+}
+
+void AnalyticSinkLocalState::_ensure_spill_batch_store() {
+ if (_batch_store) {
+ return;
+ }
+ DCHECK_EQ(_open_partition_rows, 0);
+ DCHECK(_batch_partition_ends.empty());
+ _batch_store =
std::make_unique<AnalyticSpillBatchStore>(operator_profile(),
_parent->node_id(),
+
_has_peer_group_functions);
+ _batch_partition_results.resize(_agg_functions_size);
+ _batch_function_parameters.resize(_agg_functions_size);
+ for (size_t i = 0; i < _agg_functions_size; ++i) {
+ if (_spill_strategies[i] == WindowSpillStrategy::PARTITION_REDUCE) {
+ _batch_partition_results[i] =
_agg_functions[i]->data_type()->create_column();
+ }
+ }
+}
+
+Status AnalyticSinkLocalState::_finish_spill_partition(RuntimeState* state) {
+ DCHECK_GT(_open_partition_rows, 0);
+ if (_has_peer_group_functions) {
+ RETURN_IF_ERROR(_batch_store->append_peer_group_end(state,
_open_partition_end));
+ }
+ for (size_t i = 0; i < _agg_functions_size; ++i) {
+ switch (_spill_strategies[i]) {
+ case WindowSpillStrategy::PARTITION_REDUCE: {
+ auto& result = _batch_partition_results[i];
+ if (_result_column_nullable_flags[i]) {
+ if (_use_null_result[i]) {
+ result->insert_default();
+ } else {
+ auto* nullable =
assert_cast<ColumnNullable*>(result.get());
+ nullable->get_null_map_data().push_back(0);
+ _agg_functions[i]->insert_result_info(
+ _fn_place_ptr + _offsets_of_aggregate_states[i],
+ &nullable->get_nested_column());
+ }
+ } else {
+ _agg_functions[i]->insert_result_info(
+ _fn_place_ptr + _offsets_of_aggregate_states[i],
result.get());
+ }
+ DCHECK_EQ(result->size(), _batch_partition_ends.size() + 1);
+ break;
+ }
+ case WindowSpillStrategy::PARTITION_CARDINALITY:
+
_batch_function_parameters[i].push_back(_open_partition_parameters[i]);
+ break;
+ case WindowSpillStrategy::PEER_GROUP:
+ break;
+ case WindowSpillStrategy::UNSUPPORTED:
+ DORIS_CHECK(false);
+ break;
+ }
+ }
+ _batch_partition_ends.push_back(_open_partition_end);
+ COUNTER_SET(_max_partition_rows,
+ std::max<int64_t>(_max_partition_rows->value(),
_open_partition_rows));
+ _reset_agg_status();
+ _open_partition_rows = 0;
+ _last_order_keys.clear();
+ return Status::OK();
+}
+
+Status AnalyticSinkLocalState::_seal_spill_batch(RuntimeState* state) {
+ DCHECK(_batch_store != nullptr);
+ DCHECK_EQ(_open_partition_rows, 0);
+ DCHECK(!_batch_partition_ends.empty());
+ DCHECK_EQ(_batch_partition_ends.back(), _batch_store->rows());
+
+ // Once published, an in-memory batch is owned by the source side and can
no longer be
+ // reclaimed through this sink's revoke callback. A batch that reaches the
proactive spill
+ // limit is spilled before it is published; smaller batches stay in memory.
+ if (!_batch_store->is_spilled() &&
+ std::cmp_greater_equal(_batch_store->revocable_mem_size(),
+ state->spill_analytic_sink_mem_limit_bytes())) {
+ RETURN_IF_ERROR(_spill_batch_store(state));
+ }
+
+ auto batch = std::make_shared<AnalyticSpillBatch>();
+ const bool spilled = _batch_store->is_spilled() ||
_batch_store->has_spilled_peer_groups();
+ const auto partition_count =
static_cast<int64_t>(_batch_partition_ends.size());
+ COUNTER_UPDATE(_peer_group_metadata_bytes,
_batch_store->peer_group_metadata_bytes());
+ RETURN_IF_ERROR(_batch_store->seal(state, batch.get()));
+
+ auto& parent = _parent->cast<AnalyticSinkOperatorX>();
+ batch->partition_ends = std::move(_batch_partition_ends);
+ batch->strategies = _spill_strategies;
+ batch->peer_functions = parent._window_spill_peer_functions;
+ batch->function_parameters = std::move(_batch_function_parameters);
+ batch->change_to_nullable_flags = parent._change_to_nullable_flags;
+ batch->partition_results.resize(_agg_functions_size);
+ batch->result_types.resize(_agg_functions_size);
+ for (size_t i = 0; i < _agg_functions_size; ++i) {
+ batch->result_types[i] = _agg_functions[i]->data_type();
+ if (_batch_partition_results[i]) {
+ batch->partition_results[i] =
std::move(_batch_partition_results[i]);
+ }
+ }
+
+ {
+ LockGuard lock(_shared_state->buffer_mutex);
+ _shared_state->spill_batches.push(std::move(batch));
+ }
+ _dependency->set_ready_to_read();
+ COUNTER_UPDATE(spilled ? _spilled_partitions : _in_memory_partitions,
partition_count);
+
+ _batch_store.reset();
+ _batch_partition_ends.clear();
+ _batch_partition_results.clear();
+ _batch_function_parameters.clear();
+ _update_spill_memory_usage();
+ return Status::OK();
+}
+
+Status AnalyticSinkLocalState::_sink_spill(RuntimeState* state, Block*
input_block, bool eos) {
+ RETURN_IF_CANCELLED(state);
+ if (input_block->rows() > 0) {
+ RETURN_IF_ERROR(_process_spill_block(state, input_block));
+ }
+ return _finish_spill_sink_call(state, eos);
+}
+
+Status AnalyticSinkLocalState::_finish_spill_sink_call(RuntimeState* state,
bool eos) {
+ if (eos) {
Review Comment:
[P2] Recompute admission for the empty EOS call. A preceding non-EOS call
can spill the open partition and leave a large PBlock/compression peak in
`_reserve_mem_size` even after its buffers are flushed and
`revocable_mem_size()` is zero. PipelineTask asks for that stale peak before
entering this branch; if the workload group stays above its high watermark, the
sink remains paused while the source waits for this unpublished batch, and the
pause handler can cancel the query on timeout. Size the pending seal/close work
from the current state instead of carrying the prior write peak into EOS.
--
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]