This is an automated email from the ASF dual-hosted git repository.
HappenLee pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 034a27ca41c [improvement](agg) Bound streaming pre-aggregation memory
by the query memory limit (#67140)
034a27ca41c is described below
commit 034a27ca41c1ae9d4955c204c6cc97b41f6c4b62
Author: Jerry Hu <[email protected]>
AuthorDate: Thu Oct 8 16:13:56 2026 +0800
[improvement](agg) Bound streaming pre-aggregation memory by the query
memory limit (#67140)
### What problem does this PR solve?
Issue Number: None
Problem Summary: Streaming pre-aggregation had a fixed spill-only memory
threshold, while specialized distinct pre-aggregation had no byte threshold.
Derive a shared per-task threshold from the live query memory limit and
preserve the explicit spill setting as an additional bound. Keep pass-through
rows out of pushed local LIMIT truncation, preserve TopN candidate boundaries
across pass-through and recovery, and use a block-scoped serialization arena.
### Release note
Streaming pre-aggregation now uses a query-derived memory threshold at
block boundaries, with a target floor clipped by the per-task query share. The
explicit spill threshold remains an additional bound when spilling is enabled.
Existing retained-state accounting and distinct recovery limitations remain
documented in the review threads.
### Check List (For Author)
- Test: GitHub COMPILE, BE UT, FE UT, Cloud UT, regression and performance
checks passed for the reviewed head. Local review ran git diff --check,
clang-format 16 checks and build hygiene; no local build or unit-test rerun.
- Behavior changed: Yes. Pre-aggregation can pass rows through when its
query-derived threshold is exceeded, including without spilling.
- Does this need documentation: No. No new public setting or interface.
---
.../distinct_streaming_aggregation_operator.cpp | 50 +++-
.../distinct_streaming_aggregation_operator.h | 8 +-
be/src/exec/operator/streaming_agg_memory_limit.h | 65 ++++++
.../operator/streaming_aggregation_operator.cpp | 90 ++++----
.../exec/operator/streaming_aggregation_operator.h | 14 +-
...istinct_streaming_aggregation_operator_test.cpp | 86 +++++++
.../exec/operator/streaming_agg_operator_test.cpp | 256 ++++++++++++++++++++-
.../java/org/apache/doris/qe/SessionVariable.java | 6 +-
8 files changed, 523 insertions(+), 52 deletions(-)
diff --git a/be/src/exec/operator/distinct_streaming_aggregation_operator.cpp
b/be/src/exec/operator/distinct_streaming_aggregation_operator.cpp
index 8d751c5b9d0..c9b61f1f832 100644
--- a/be/src/exec/operator/distinct_streaming_aggregation_operator.cpp
+++ b/be/src/exec/operator/distinct_streaming_aggregation_operator.cpp
@@ -23,8 +23,12 @@
#include <utility>
#include "common/compiler_util.h" // IWYU pragma: keep
+#include "exec/operator/streaming_agg_memory_limit.h"
#include "exec/operator/streaming_agg_min_reduction.h"
#include "exprs/vectorized_agg_fn.h"
+#include "runtime/query_context.h"
+#include "runtime/workload_management/memory_context.h"
+#include "runtime/workload_management/resource_context.h"
namespace doris {
class ExecNode;
@@ -52,6 +56,7 @@ Status DistinctStreamingAggLocalState::init(RuntimeState*
state, LocalStateInfo&
_hash_table_input_counter =
ADD_COUNTER(Base::custom_profile(), "HashTableInputCount",
TUnit::UNIT);
_hash_table_size_counter = ADD_COUNTER(custom_profile(), "HashTableSize",
TUnit::UNIT);
+ _memory_use_limit = ADD_COUNTER(custom_profile(), "MemoryUseLimit",
TUnit::BYTES);
_insert_keys_to_column_timer = ADD_TIMER(custom_profile(),
"InsertKeysToColumnTime");
return Status::OK();
@@ -67,9 +72,23 @@ Status DistinctStreamingAggLocalState::open(RuntimeState*
state) {
RETURN_IF_ERROR(p._probe_expr_ctxs[i]->clone(state,
_probe_expr_ctxs[i]));
}
RETURN_IF_ERROR(_init_hash_method(_probe_expr_ctxs));
+ COUNTER_SET(_memory_use_limit,
static_cast<int64_t>(p._memory_limit(state)));
return Status::OK();
}
+size_t DistinctStreamingAggLocalState::_memory_usage() const {
+ size_t usage = _arena.size();
+ std::visit(Overload {[&](std::monostate& arg) -> void {
+ throw doris::Exception(ErrorCode::INTERNAL_ERROR,
+ "uninited hash table");
+ },
+ [&](auto& agg_method) {
+ usage +=
agg_method.hash_table->get_buffer_size_in_bytes();
+ }},
+ _agg_data->method_variant);
+ return usage;
+}
+
bool DistinctStreamingAggLocalState::_should_expand_preagg_hash_tables() {
if (!_should_expand_hash_table) {
return false;
@@ -176,8 +195,15 @@ Status
DistinctStreamingAggLocalState::_distinct_pre_agg_with_serialized_key(
const uint32_t rows = (uint32_t)in_block->rows();
_distinct_row.clear();
- if (_parent->cast<DistinctStreamingAggOperatorX>()._is_streaming_preagg &&
low_memory_mode()) {
- _stop_emplace_flag = true;
+ auto& parent = _parent->cast<DistinctStreamingAggOperatorX>();
+ if (parent._is_streaming_preagg) {
+ const auto memory_limit = parent._memory_limit(state());
+ COUNTER_SET(_memory_use_limit, static_cast<int64_t>(memory_limit));
+ // Latching is safe under a pushed-down LIMIT as well: push() stops
truncating raw rows
+ // once the flag is set and the global stage applies the limit again.
+ if (low_memory_mode() || (memory_limit > 0 && _memory_usage() >
memory_limit)) {
+ _stop_emplace_flag = true;
+ }
}
if (!_stop_emplace_flag) {
@@ -327,6 +353,9 @@
DistinctStreamingAggOperatorX::DistinctStreamingAggOperatorX(ObjectPool* pool, i
Status DistinctStreamingAggOperatorX::init(const TPlanNode& tnode,
RuntimeState* state) {
RETURN_IF_ERROR(StatefulOperatorX<DistinctStreamingAggLocalState>::init(tnode,
state));
+ _spill_streaming_agg_mem_limit =
state->query_options().__isset.spill_streaming_agg_mem_limit
+ ?
state->query_options().spill_streaming_agg_mem_limit
+ : 0;
// ignore return status for now , so we need to introduce ExecNode::init()
RETURN_IF_ERROR(VExpr::create_expr_trees(tnode.agg_node.grouping_exprs,
_probe_expr_ctxs));
@@ -343,6 +372,18 @@ Status
DistinctStreamingAggOperatorX::prepare(RuntimeState* state) {
return Status::OK();
}
+size_t DistinctStreamingAggOperatorX::_memory_limit(RuntimeState* state) const
{
+ if (!_is_streaming_preagg) {
+ return 0;
+ }
+ // Same rule as StreamingAggOperatorX: the fixed bound only applies when
spilling is enabled.
+ const int64_t fixed_limit =
+ state->enable_spill() ?
static_cast<int64_t>(_spill_streaming_agg_mem_limit) : 0;
+ return streaming_agg_memory_limit(
+
state->get_query_ctx()->resource_ctx()->memory_context()->mem_limit(),
parallel_tasks(),
+ fixed_limit);
+}
+
void DistinctStreamingAggOperatorX::init_make_nullable(RuntimeState* state) {
_output_tuple_desc =
state->desc_tbl().get_tuple_descriptor(_output_tuple_id);
@@ -366,7 +407,10 @@ Status DistinctStreamingAggOperatorX::push(RuntimeState*
state, Block* in_block,
RETURN_IF_ERROR(local_state._distinct_pre_agg_with_serialized_key(
in_block, local_state._aggregated_block.get()));
// Prevents exceeding the row limit when the aggregated block reaches or
equals the threshold.
- if (_limit != -1 &&
+ // Pass-through rows are not deduplicated; counting them against a
pushed-down limit would let
+ // duplicates consume the allowance and stop the child before enough
distinct keys reached the
+ // global stage (which applies the limit again).
+ if (_limit != -1 && !local_state._stop_emplace_flag &&
(local_state._num_rows_returned +
local_state._aggregated_block->rows()) >= _limit) {
auto limit_rows = _limit - local_state._num_rows_returned;
local_state._aggregated_block->set_num_rows(limit_rows);
diff --git a/be/src/exec/operator/distinct_streaming_aggregation_operator.h
b/be/src/exec/operator/distinct_streaming_aggregation_operator.h
index 3d5304aff1e..76b55542025 100644
--- a/be/src/exec/operator/distinct_streaming_aggregation_operator.h
+++ b/be/src/exec/operator/distinct_streaming_aggregation_operator.h
@@ -55,6 +55,7 @@ private:
ColumnRawPtrs& key_columns,
const uint32_t num_rows);
void _make_nullable_output_key(Block* block);
bool _should_expand_preagg_hash_tables();
+ size_t _memory_usage() const;
void _swap_cache_block(Block* block) {
DCHECK(!_cache_block.is_empty_column());
@@ -82,6 +83,7 @@ private:
RuntimeProfile::Counter* _hash_table_emplace_timer = nullptr;
RuntimeProfile::Counter* _hash_table_input_counter = nullptr;
RuntimeProfile::Counter* _hash_table_size_counter = nullptr;
+ RuntimeProfile::Counter* _memory_use_limit = nullptr;
RuntimeProfile::Counter* _insert_keys_to_column_timer = nullptr;
bool _is_single_backend = false;
@@ -145,7 +147,7 @@ public:
private:
friend class DistinctStreamingAggLocalState;
-
+ size_t _memory_limit(RuntimeState* state) const;
void init_make_nullable(RuntimeState* state);
TupleId _output_tuple_id;
TupleDescriptor* _output_tuple_desc = nullptr;
@@ -159,6 +161,10 @@ private:
// If _is_streaming_preagg = true, deduplication will be abandoned in
cases where the deduplication rate is low.
bool _is_streaming_preagg = false;
+ /// When spilling is enabled, the pre-agg should not occupy too much
memory: session variable
+ /// `spill_streaming_agg_mem_limit` (0 = none), combined with the
query-limit-based budget in
+ /// `_memory_limit()`.
+ size_t _spill_streaming_agg_mem_limit = 0;
};
/// Instantiated once in operator.cpp; suppresses per-TU implicit
instantiation.
diff --git a/be/src/exec/operator/streaming_agg_memory_limit.h
b/be/src/exec/operator/streaming_agg_memory_limit.h
new file mode 100644
index 00000000000..9d6b9c15941
--- /dev/null
+++ b/be/src/exec/operator/streaming_agg_memory_limit.h
@@ -0,0 +1,65 @@
+// 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.
+
+#pragma once
+
+#include <algorithm>
+#include <cstddef>
+#include <cstdint>
+
+#include "exec/operator/streaming_agg_min_reduction.h"
+
+namespace doris {
+
+// Per-task memory budget of a streaming pre-aggregation (hash table + arenas).
+//
+// 1. The budget is one fifth of the current query memory limit, shared evenly
among the
+// `parallel_tasks` instances of the operator on this BE. The query limit
is read on every
+// call, so a limit lowered or restored by the workload group manager takes
effect at once.
+// 2. It never drops below `min_memory_limit`, which is the last cache tier of
the min-reduction
+// table doubled (the budget also counts the key/state arenas, not only the
bucket array), so a
+// small query limit does not disable pre-aggregation altogether. The floor
itself is capped by
+// the per-task share of the query limit, so the pre-aggregation alone can
never exceed it.
+// 3. `fixed_limit` is an explicit upper bound on top of that; 0 means "no
fixed bound". Callers
+// pass the session variable `spill_streaming_agg_mem_limit` when spilling
is enabled (the
+// downstream agg can spill, the pre-agg cannot, so it must stay small) and
0 otherwise. It is
+// applied last so that a user who sets it explicitly always gets what they
asked for.
+//
+// Returns 0 when neither limit is known, which callers treat as "no cap".
+inline size_t streaming_agg_memory_limit(int64_t query_memory_limit, int
parallel_tasks,
+ int64_t fixed_limit) {
+ if (query_memory_limit <= 0) {
+ return fixed_limit > 0 ? static_cast<size_t>(fixed_limit) : 0;
+ }
+
+ constexpr int64_t memory_limit_divisor = 5;
+ constexpr int64_t min_memory_limit =
+ 2LL * STREAMING_HT_MIN_REDUCTION[STREAMING_HT_MIN_REDUCTION_SIZE -
1].min_ht_mem;
+
+ // A known positive query limit must never collapse to the "no cap"
sentinel 0, even when it is
+ // smaller than the number of tasks.
+ const int64_t per_task_query_limit =
+ std::max<int64_t>(query_memory_limit / std::max(parallel_tasks,
1), 1);
+ int64_t limit = per_task_query_limit / memory_limit_divisor;
+ limit = std::max(limit, std::min(min_memory_limit, per_task_query_limit));
+ if (fixed_limit > 0) {
+ limit = std::min(limit, fixed_limit);
+ }
+ return static_cast<size_t>(limit);
+}
+
+} // namespace doris
diff --git a/be/src/exec/operator/streaming_aggregation_operator.cpp
b/be/src/exec/operator/streaming_aggregation_operator.cpp
index de35cec0da7..61637187817 100644
--- a/be/src/exec/operator/streaming_aggregation_operator.cpp
+++ b/be/src/exec/operator/streaming_aggregation_operator.cpp
@@ -26,11 +26,16 @@
#include "common/compiler_util.h" // IWYU pragma: keep
#include "core/column/column_fixed_length_object.h"
#include "exec/operator/operator.h"
+#include "exec/operator/streaming_agg_memory_limit.h"
#include "exec/operator/streaming_agg_min_reduction.h"
#include "exprs/aggregate/aggregate_function_count.h"
#include "exprs/aggregate/aggregate_function_simple_factory.h"
#include "exprs/vectorized_agg_fn.h"
#include "exprs/vslot_ref.h"
+#include "runtime/query_context.h"
+#include "runtime/workload_management/io_context.h"
+#include "runtime/workload_management/memory_context.h"
+#include "runtime/workload_management/resource_context.h"
namespace doris {
class RuntimeState;
@@ -66,6 +71,7 @@ Status StreamingAggLocalState::init(RuntimeState* state,
LocalStateInfo& info) {
_hash_table_input_counter =
ADD_COUNTER(Base::custom_profile(), "HashTableInputCount",
TUnit::UNIT);
_hash_table_size_counter = ADD_COUNTER(custom_profile(), "HashTableSize",
TUnit::UNIT);
+ _memory_use_limit = ADD_COUNTER(custom_profile(), "MemoryUseLimit",
TUnit::BYTES);
_streaming_agg_timer = ADD_TIMER(custom_profile(), "StreamingAggTime");
_build_timer = ADD_TIMER(custom_profile(), "BuildTime");
_expr_timer = ADD_TIMER(Base::custom_profile(), "ExprTime");
@@ -134,6 +140,8 @@ Status StreamingAggLocalState::open(RuntimeState* state) {
null_directions = p._null_directions;
order_directions = p._order_directions;
+ COUNTER_SET(_memory_use_limit,
static_cast<int64_t>(p._memory_limit(state)));
+
return Status::OK();
}
@@ -289,9 +297,9 @@ bool StreamingAggLocalState::_should_not_do_pre_agg(size_t
rows) {
// But for fixed hash map, it never need to expand
auto& p = Base::_parent->template cast<StreamingAggOperatorX>();
bool ret_flag = false;
- const auto spill_streaming_agg_mem_limit =
p._spill_streaming_agg_mem_limit;
- const bool used_too_much_memory =
- spill_streaming_agg_mem_limit > 0 && _memory_usage() >
spill_streaming_agg_mem_limit;
+ const auto memory_limit = p._memory_limit(state());
+ COUNTER_SET(_memory_use_limit, static_cast<int64_t>(memory_limit));
+ const bool used_too_much_memory = memory_limit > 0 && _memory_usage() >
memory_limit;
std::visit(
Overload {
[&](std::monostate& arg) {
@@ -342,22 +350,25 @@ Status
StreamingAggLocalState::_pre_agg_with_serialized_key(doris::Block* in_blo
_places.resize(rows);
if (_should_not_do_pre_agg(rows)) {
+ // Serializing a row as a single-row state may allocate from the arena
(collect, map,
+ // foreach, ...). Those bytes are dead once the state is serialized,
so they must not
+ // accumulate in the operator-lifetime `_agg_arena_pool` for every
pass-through block.
+ Arena pass_through_arena;
if (limit > 0) {
DCHECK(do_sort_limit);
- if (need_do_sort_limit == -1) {
- const size_t hash_table_size = _get_hash_table_size();
- need_do_sort_limit = hash_table_size >= limit ? 1 : 0;
- if (need_do_sort_limit == 1) {
- build_limit_heap(hash_table_size);
- }
- }
-
+ // The heap is only ever seeded on the aggregation path: the hash
table can only
+ // grow there, and that path latches as soon as it holds `limit`
groups. So here
+ // the heap either exists already or the table is still too small
to build one.
if (need_do_sort_limit == 1) {
+ // Only drop rows that the current boundary already excludes.
The boundary must
+ // not be tightened from pass-through rows: they never enter
the hash table, so
+ // the heap cannot tell a new group from one it already holds
and would count
+ // the same group twice, discarding a group that belongs in
the TopN result.
+ // Aggregation resumed later refreshes the heap through hash
table insertion.
if (_do_limit_filter(rows, key_columns)) {
bool need_filter = std::find(need_computes.begin(),
need_computes.end(), 1) !=
need_computes.end();
if (need_filter) {
- _add_limit_heap_top(key_columns, rows);
Block::filter_block_internal(in_block, need_computes);
rows = (uint32_t)in_block->rows();
} else {
@@ -374,7 +385,7 @@ Status
StreamingAggLocalState::_pre_agg_with_serialized_key(doris::Block* in_blo
for (int i = 0; i != _aggregate_evaluators.size(); ++i) {
SCOPED_TIMER(_insert_values_to_column_timer);
RETURN_IF_ERROR(_aggregate_evaluators[i]->streaming_agg_serialize_to_column(
- in_block, columns[i + key_size], rows,
_agg_arena_pool));
+ in_block, columns[i + key_size], rows,
pass_through_arena));
}
for (int i = 0; i < key_size; ++i) {
columns[i]->insert_range_from(*key_columns[i], 0, rows);
@@ -393,7 +404,7 @@ Status
StreamingAggLocalState::_pre_agg_with_serialized_key(doris::Block* in_blo
for (int i = 0; i != _aggregate_evaluators.size(); ++i) {
SCOPED_TIMER(_insert_values_to_column_timer);
RETURN_IF_ERROR(_aggregate_evaluators[i]->streaming_agg_serialize_to_column(
- in_block, value_columns[i], rows, _agg_arena_pool));
+ in_block, value_columns[i], rows, pass_through_arena));
}
ColumnsWithTypeAndName columns_with_schema;
@@ -692,21 +703,6 @@ void StreamingAggLocalState::build_limit_heap(size_t
hash_table_size) {
limit_columns_min = limit_heap.top()._row_id;
}
-void StreamingAggLocalState::_add_limit_heap_top(ColumnRawPtrs& key_columns,
size_t rows) {
- for (int i = 0; i < rows; ++i) {
- if (cmp_res[i] == 1 && need_computes[i]) {
- for (int j = 0; j < key_columns.size(); ++j) {
- limit_columns[j]->insert_from(*key_columns[j], i);
- }
- limit_heap.emplace(limit_columns[0]->size() - 1, limit_columns,
order_directions,
- null_directions);
- limit_heap.pop();
- limit_columns_min = limit_heap.top()._row_id;
- break;
- }
- }
-}
-
void StreamingAggLocalState::_refresh_limit_heap(size_t i, ColumnRawPtrs&
key_columns) {
for (int j = 0; j < key_columns.size(); ++j) {
limit_columns[j]->insert_from(*key_columns[j], i);
@@ -954,15 +950,9 @@ Status StreamingAggOperatorX::init(const TPlanNode& tnode,
RuntimeState* state)
_aggregate_evaluators.push_back(evaluator);
}
- if (state->enable_spill()) {
- // If spill enabled, the streaming agg should not occupy too much
memory.
- _spill_streaming_agg_mem_limit =
- state->query_options().__isset.spill_streaming_agg_mem_limit
- ? state->query_options().spill_streaming_agg_mem_limit
- : 0;
- } else {
- _spill_streaming_agg_mem_limit = 0;
- }
+ _spill_streaming_agg_mem_limit =
state->query_options().__isset.spill_streaming_agg_mem_limit
+ ?
state->query_options().spill_streaming_agg_mem_limit
+ : 0;
const auto& agg_functions = tnode.agg_node.aggregate_functions;
auto is_merge = std::any_of(agg_functions.cbegin(), agg_functions.cend(),
@@ -996,6 +986,22 @@ Status StreamingAggOperatorX::init(const TPlanNode& tnode,
RuntimeState* state)
return Status::OK();
}
+size_t StreamingAggOperatorX::_memory_limit(RuntimeState* state) const {
+ // When spilling is enabled, the streaming agg should not occupy too much
memory: the downstream
+ // agg can spill, the pre-agg cannot.
+ const int64_t fixed_limit =
+ state->enable_spill() ?
static_cast<int64_t>(_spill_streaming_agg_mem_limit) : 0;
+ const size_t limit = streaming_agg_memory_limit(
+
state->get_query_ctx()->resource_ctx()->memory_context()->mem_limit(),
parallel_tasks(),
+ fixed_limit);
+ if (_low_memory_mode.load(std::memory_order_relaxed)) {
+ // Low-memory mode only ever tightens the budget.
+ constexpr size_t low_memory_mode_limit = 1024 * 1024;
+ return limit > 0 ? std::min(limit, low_memory_mode_limit) :
low_memory_mode_limit;
+ }
+ return limit;
+}
+
Status StreamingAggOperatorX::prepare(RuntimeState* state) {
RETURN_IF_ERROR(StatefulOperatorX<StreamingAggLocalState>::prepare(state));
@@ -1106,14 +1112,20 @@ Status StreamingAggOperatorX::pull(RuntimeState* state,
Block* block, bool* eos)
auto& local_state = get_local_state(state);
SCOPED_PEAK_MEM(&local_state._estimate_memory_usage);
if (!local_state._pre_aggregated_block->empty()) {
+ // Pass-through rows are neither aggregated nor deduplicated, so a
limit pushed down to
+ // this local stage must not count them: duplicates would consume the
allowance and stop
+ // the child before enough distinct keys reached the global stage,
which applies the
+ // limit again on the final result. They are still processed rows for
the query
+ // statistics, which reached_limit() would otherwise have accounted
for.
local_state._pre_aggregated_block->swap(*block);
+
state->get_query_ctx()->resource_ctx()->io_context()->update_process_rows(block->rows());
} else {
RETURN_IF_ERROR(local_state._get_results_with_serialized_key(state,
block, eos));
local_state.make_nullable_output_key(block);
// dispose the having clause, should not be execute in prestreaming agg
RETURN_IF_ERROR(local_state.filter_block(local_state._conjuncts,
block));
+ local_state.reached_limit(block, eos);
}
- local_state.reached_limit(block, eos);
return Status::OK();
}
diff --git a/be/src/exec/operator/streaming_aggregation_operator.h
b/be/src/exec/operator/streaming_aggregation_operator.h
index 48fae130213..e5394a8db70 100644
--- a/be/src/exec/operator/streaming_aggregation_operator.h
+++ b/be/src/exec/operator/streaming_aggregation_operator.h
@@ -19,6 +19,7 @@
#include <stdint.h>
+#include <atomic>
#include <memory>
#include "common/status.h"
@@ -53,7 +54,6 @@ private:
friend class StatefulOperatorX;
size_t _memory_usage() const;
- void _add_limit_heap_top(ColumnRawPtrs& key_columns, size_t rows);
bool _do_limit_filter(size_t num_rows, ColumnRawPtrs& key_columns);
void _refresh_limit_heap(size_t i, ColumnRawPtrs& key_columns);
@@ -87,6 +87,7 @@ private:
RuntimeProfile::Counter* _hash_table_memory_usage = nullptr;
RuntimeProfile::HighWaterMarkCounter* _serialize_key_arena_memory_usage =
nullptr;
RuntimeProfile::Counter* _hash_table_size_counter = nullptr;
+ RuntimeProfile::Counter* _memory_use_limit = nullptr;
RuntimeProfile::Counter* _get_results_timer = nullptr;
RuntimeProfile::Counter* _hash_table_iterate_timer = nullptr;
RuntimeProfile::Counter* _insert_keys_to_column_timer = nullptr;
@@ -220,7 +221,7 @@ public:
Status push(RuntimeState* state, Block* input_block, bool eos) const
override;
bool need_more_input_data(RuntimeState* state) const override;
void set_low_memory_mode(RuntimeState* state) override {
- _spill_streaming_agg_mem_limit = 1024 * 1024;
+ _low_memory_mode.store(true, std::memory_order_relaxed);
}
DataDistribution required_data_distribution(RuntimeState* state) const
override {
if (_child && _child->is_hash_join_probe() &&
@@ -245,6 +246,8 @@ public:
private:
friend class StreamingAggLocalState;
+ size_t _memory_limit(RuntimeState* state) const;
+
MOCK_FUNCTION Status _init_probe_expr_ctx(RuntimeState* state);
MOCK_FUNCTION Status _init_aggregate_evaluators(RuntimeState* state);
@@ -265,8 +268,11 @@ private:
/// The total size of the row from the aggregate functions.
size_t _total_size_of_aggregate_states = 0;
- /// When spilling is enabled, the streaming agg should not occupy too much
memory.
- size_t _spill_streaming_agg_mem_limit;
+ /// When spilling is enabled, the streaming agg should not occupy too much
memory: session
+ /// variable `spill_streaming_agg_mem_limit` (0 = none), combined with the
query-limit-based
+ /// budget in `_memory_limit()`.
+ size_t _spill_streaming_agg_mem_limit = 0;
+ std::atomic_bool _low_memory_mode = false;
// group by k1,k2
VExprContextSPtrs _probe_expr_ctxs;
std::vector<AggFnEvaluator*> _aggregate_evaluators;
diff --git
a/be/test/exec/operator/distinct_streaming_aggregation_operator_test.cpp
b/be/test/exec/operator/distinct_streaming_aggregation_operator_test.cpp
index 1e6976e07bc..7c3e732925f 100644
--- a/be/test/exec/operator/distinct_streaming_aggregation_operator_test.cpp
+++ b/be/test/exec/operator/distinct_streaming_aggregation_operator_test.cpp
@@ -25,6 +25,7 @@
#include "exec/exchange/local_exchange_source_operator.h"
#include "exec/operator/mock_operator.h"
#include "exec/operator/operator_helper.h"
+#include "runtime/workload_management/memory_context.h"
#include "testutil/column_helper.h"
#include "testutil/mock/mock_descriptors.h"
#include "testutil/mock/mock_slot_ref.h"
@@ -209,4 +210,89 @@ TEST_F(DistinctStreamingAggOperatorTest, test3) {
{ EXPECT_TRUE(op->close(state.get())); }
}
+TEST_F(DistinctStreamingAggOperatorTest, refresh_memory_limit) {
+ op->_is_streaming_preagg = true;
+ op->set_parallel_tasks(2);
+ auto* memory_context =
state->get_query_ctx()->resource_ctx()->memory_context();
+ memory_context->set_mem_limit(1024LL * 1024 * 1024);
+ create_op({std::make_shared<DataTypeInt64>()},
{std::make_shared<DataTypeInt64>()});
+
+ // 1GB / 2 tasks / 5
+ auto* memory_use_limit =
local_state->custom_profile()->get_counter("MemoryUseLimit");
+ ASSERT_NE(memory_use_limit, nullptr);
+ EXPECT_EQ(memory_use_limit->value(), 1024LL * 1024 * 1024 / 2 / 5);
+
+ auto block = ColumnHelper::create_block<DataTypeInt64>({1, 2, 3, 4});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_FALSE(local_state->_stop_emplace_flag);
+
+ // spill_streaming_agg_mem_limit is ignored while spilling is disabled ...
+ op->_spill_streaming_agg_mem_limit = 4 * 1024 * 1024;
+ block = ColumnHelper::create_block<DataTypeInt64>({5});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_FALSE(local_state->_stop_emplace_flag);
+ EXPECT_EQ(memory_use_limit->value(), 1024LL * 1024 * 1024 / 2 / 5);
+
+ // ... and caps the distinct pre-agg once spilling is enabled.
+ state->set_enable_spill(true);
+ block = ColumnHelper::create_block<DataTypeInt64>({6});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_FALSE(local_state->_stop_emplace_flag);
+ EXPECT_EQ(memory_use_limit->value(), 4 * 1024 * 1024);
+
+ // A tiny query limit: the floor is capped by the per-task share (10 / 2 =
5 bytes), the
+ // hash table already exceeds it, so the pre-agg gives up and passes rows
through.
+ memory_context->set_mem_limit(10);
+ block = ColumnHelper::create_block<DataTypeInt64>({1, 1});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_TRUE(local_state->_stop_emplace_flag);
+ EXPECT_EQ(memory_use_limit->value(), 5);
+ EXPECT_EQ(local_state->_aggregated_block->rows(), 8);
+}
+
+TEST_F(DistinctStreamingAggOperatorTest, pushed_limit_with_memory_limit) {
+ op->_is_streaming_preagg = true;
+ op->_limit = 2;
+ op->set_parallel_tasks(2);
+ auto* memory_context =
state->get_query_ctx()->resource_ctx()->memory_context();
+ memory_context->set_mem_limit(1024LL * 1024 * 1024);
+ create_op({std::make_shared<DataTypeInt64>()},
{std::make_shared<DataTypeInt64>()});
+
+ // Within the budget a pushed-down LIMIT works as usual: duplicates are
removed and the
+ // operator stops once `limit` distinct keys are out.
+ auto block = ColumnHelper::create_block<DataTypeInt64>({1, 1});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_FALSE(local_state->_stop_emplace_flag);
+ EXPECT_FALSE(local_state->_reach_limit);
+ EXPECT_EQ(local_state->_aggregated_block->rows(), 1);
+
+ // Budget exceeded: the operator latches into pass-through and no longer
truncates, so the
+ // global stage still sees every key that may be distinct (the limit is
applied there).
+ memory_context->set_mem_limit(10);
+ block = ColumnHelper::create_block<DataTypeInt64>({1, 1, 2, 3});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_TRUE(local_state->_stop_emplace_flag);
+ EXPECT_FALSE(local_state->_reach_limit);
+ EXPECT_EQ(local_state->_aggregated_block->rows(), 5);
+}
+
+TEST_F(DistinctStreamingAggOperatorTest,
pass_through_does_not_consume_pushed_limit) {
+ op->_is_streaming_preagg = true;
+ op->_limit = 2;
+ create_op({std::make_shared<DataTypeInt64>()},
{std::make_shared<DataTypeInt64>()});
+
+ auto block = ColumnHelper::create_block<DataTypeInt64>({1});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_EQ(local_state->_aggregated_block->rows(), 1);
+
+ // Once the operator has permanently stopped deduplicating (low reduction
rate / low-memory
+ // mode), raw rows must not be truncated against the limit: the global
stage needs every
+ // key that may still be distinct.
+ local_state->_stop_emplace_flag = true;
+ block = ColumnHelper::create_block<DataTypeInt64>({1, 1, 2});
+ EXPECT_TRUE(op->push(state.get(), &block, false));
+ EXPECT_FALSE(local_state->_reach_limit);
+ EXPECT_EQ(local_state->_aggregated_block->rows(), 4);
+}
+
} // namespace doris
diff --git a/be/test/exec/operator/streaming_agg_operator_test.cpp
b/be/test/exec/operator/streaming_agg_operator_test.cpp
index 7e2fbb6a20b..5d2bc376321 100644
--- a/be/test/exec/operator/streaming_agg_operator_test.cpp
+++ b/be/test/exec/operator/streaming_agg_operator_test.cpp
@@ -28,7 +28,10 @@
#include "exec/operator/aggregation_source_operator.h"
#include "exec/operator/mock_operator.h"
#include "exec/operator/operator_helper.h"
+#include "exec/operator/streaming_agg_memory_limit.h"
#include "exec/operator/streaming_aggregation_operator.h"
+#include "runtime/workload_management/io_context.h"
+#include "runtime/workload_management/memory_context.h"
#include "testutil/column_helper.h"
#include "testutil/mock/mock_agg_fn_evaluator.h"
#include "testutil/mock/mock_runtime_state.h"
@@ -52,12 +55,12 @@ struct MockStreamingAggLocalState : public
StreamingAggLocalState {
bool _should_not_do_pre_agg(size_t rows) override {
static_cast<void>(_should_expand_preagg_hash_tables()); // mock the
function
static_cast<void>(_memory_usage()); // mock the
function
- static_cast<void>(
- StreamingAggLocalState::_should_not_do_pre_agg(rows)); // mock
the function
- return should_not_do_pre_agg;
+ const bool real_decision =
StreamingAggLocalState::_should_not_do_pre_agg(rows);
+ return use_real_decision ? real_decision : should_not_do_pre_agg;
}
bool should_not_do_pre_agg = false;
+ bool use_real_decision = false;
};
class MockStreamingAggOperatorChildOperator : public OperatorXBase {
@@ -104,7 +107,35 @@ struct StreamingAggOperatorTest : public testing::Test {
ObjectPool pool;
};
+TEST(StreamingAggMemoryLimitTest, budget_floor_and_fixed_limit) {
+ constexpr int64_t MB = 1024 * 1024;
+ // One fifth of the per-task share of the query limit.
+ EXPECT_EQ(streaming_agg_memory_limit(5000 * MB, 5, 0), size_t(200 * MB));
+ // parallel_tasks <= 0 is treated as 1.
+ EXPECT_EQ(streaming_agg_memory_limit(1000 * MB, 0, 0), size_t(200 * MB));
+ // Floor: a small query limit still leaves twice the last cache tier
(32MB) ...
+ EXPECT_EQ(streaming_agg_memory_limit(2048 * MB, 16, 0), size_t(32 * MB));
+ // ... but never more than the per-task share of the query limit.
+ EXPECT_EQ(streaming_agg_memory_limit(64 * MB, 16, 0), size_t(4 * MB));
+ EXPECT_EQ(streaming_agg_memory_limit(10, 2, 0), size_t(5));
+ // A positive query limit smaller than the task count still yields a cap,
never "no cap".
+ EXPECT_EQ(streaming_agg_memory_limit(1, 2, 0), size_t(1));
+ EXPECT_EQ(streaming_agg_memory_limit(1, 2, 256 * MB), size_t(1));
+ // The fixed bound (spill_streaming_agg_mem_limit) is applied last: it
never raises the
+ // budget, and an explicit small value beats the floor.
+ EXPECT_EQ(streaming_agg_memory_limit(5000 * MB, 5, 256 * MB), size_t(200 *
MB));
+ EXPECT_EQ(streaming_agg_memory_limit(50000 * MB, 5, 256 * MB), size_t(256
* MB));
+ EXPECT_EQ(streaming_agg_memory_limit(2048 * MB, 16, 8 * MB), size_t(8 *
MB));
+ // Unknown query limit: only the fixed bound, or no cap at all.
+ EXPECT_EQ(streaming_agg_memory_limit(0, 4, 256 * MB), size_t(256 * MB));
+ EXPECT_EQ(streaming_agg_memory_limit(-1, 4, 0), size_t(0));
+}
+
TEST_F(StreamingAggOperatorTest, test1) {
+ auto* memory_context =
state->get_query_ctx()->resource_ctx()->memory_context();
+ memory_context->set_mem_limit(5000LL * 1024 * 1024);
+ op->set_parallel_tasks(5);
+
op->_aggregate_evaluators.push_back(create_mock_agg_fn_evaluator(
pool, MockSlotRef::create_mock_contexts(1,
std::make_shared<DataTypeInt64>()), false,
false));
@@ -133,9 +164,13 @@ TEST_F(StreamingAggOperatorTest, test1) {
local_state =
static_cast<MockStreamingAggLocalState*>(state->get_local_state(op->operator_id()));
EXPECT_TRUE(local_state->open(state.get()).ok());
+ auto* memory_use_limit =
local_state->custom_profile()->get_counter("MemoryUseLimit");
+ ASSERT_NE(memory_use_limit, nullptr);
+ EXPECT_EQ(memory_use_limit->value(), 200 * 1024 * 1024);
}
{
+ memory_context->set_mem_limit(2500LL * 1024 * 1024);
Block block {
ColumnHelper::create_column_with_name<DataTypeInt64>({1, 1, 2,
2, 2, 3}),
ColumnHelper::create_column_with_name<DataTypeInt64>({1, 1,
100, 100, 100, 1000})};
@@ -144,9 +179,28 @@ TEST_F(StreamingAggOperatorTest, test1) {
EXPECT_EQ(local_state->_get_hash_table_size(), 3);
EXPECT_TRUE(op->need_more_input_data(state.get()));
+
EXPECT_EQ(local_state->custom_profile()->get_counter("MemoryUseLimit")->value(),
+ 100 * 1024 * 1024);
+ }
+
+ {
+ // With spilling enabled, spill_streaming_agg_mem_limit caps the
budget.
+ state->set_enable_spill(true);
+ op->_spill_streaming_agg_mem_limit = 16 * 1024 * 1024;
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1,
2, 3}),
+ ColumnHelper::create_column_with_name<DataTypeInt64>({1,
100, 1000})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+
+ EXPECT_EQ(local_state->_get_hash_table_size(), 3);
+
EXPECT_EQ(local_state->custom_profile()->get_counter("MemoryUseLimit")->value(),
+ 16 * 1024 * 1024);
}
{
+ // Low-memory mode never raises a tighter bound.
+ op->_spill_streaming_agg_mem_limit = 512 * 1024;
+ op->set_low_memory_mode(state.get());
Block block {
ColumnHelper::create_column_with_name<DataTypeInt64>({2, 2, 2,
2, 4, 4}),
ColumnHelper::create_column_with_name<DataTypeInt64>({1, 1,
100, 100, 100, 1000})};
@@ -155,6 +209,202 @@ TEST_F(StreamingAggOperatorTest, test1) {
EXPECT_EQ(local_state->_get_hash_table_size(), 4);
EXPECT_TRUE(op->need_more_input_data(state.get()));
+
EXPECT_EQ(local_state->custom_profile()->get_counter("MemoryUseLimit")->value(),
+ 512 * 1024);
+ }
+
+ { EXPECT_TRUE(local_state->close(state.get()).ok()); }
+}
+
+TEST_F(StreamingAggOperatorTest, memory_limit_pass_through_and_recover) {
+ // A real aggregate function so that the pass-through serialization path
is exercised.
+ op->_aggregate_evaluators.push_back(create_agg_fn(pool, "sum",
+
{std::make_shared<DataTypeInt64>()},
+
std::make_shared<DataTypeInt64>(), false));
+ op->_pool = &pool;
+ op->_needs_finalize = false;
+ // A LIMIT pushed down to this local stage (pure distinct through the
regular operator).
+ op->_limit = 2;
+ op->set_parallel_tasks(5);
+
+ EXPECT_TRUE(op->set_child(child_op));
+ EXPECT_TRUE(op->prepare(state.get()).ok());
+ op->_probe_expr_ctxs = MockSlotRef::create_mock_contexts(1,
std::make_shared<DataTypeInt64>());
+
+ {
+ auto local_state =
std::make_unique<MockStreamingAggLocalState>(state.get(), op.get());
+ LocalStateInfo info {.parent_profile = &profile,
+ .scan_ranges = {},
+ .shared_state = nullptr,
+ .shared_state_map = {},
+ .task_idx = 0};
+ EXPECT_TRUE(local_state->init(state.get(), info).ok());
+ state->resize_op_id_to_local_state(-100);
+ state->emplace_local_state(op->operator_id(), std::move(local_state));
+ }
+ local_state =
+
static_cast<MockStreamingAggLocalState*>(state->get_local_state(op->operator_id()));
+ EXPECT_TRUE(local_state->open(state.get()).ok());
+ // Let the production decision drive push() instead of the mock's fixed
answer.
+ local_state->use_real_decision = true;
+ auto* memory_context =
state->get_query_ctx()->resource_ctx()->memory_context();
+ auto* memory_use_limit =
local_state->custom_profile()->get_counter("MemoryUseLimit");
+ auto* io_context = state->get_query_ctx()->resource_ctx()->io_context();
+
+ {
+ memory_context->set_mem_limit(5000LL * 1024 * 1024);
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1,
1, 100}),
+ ColumnHelper::create_column_with_name<DataTypeInt64>({1,
1, 2})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 2);
+ EXPECT_EQ(local_state->_pre_aggregated_block->rows(), 0);
+ }
+
+ {
+ // Query limit lowered below the current usage: the block passes
through untouched, the
+ // hash table does not grow, and the duplicate keys do not count
against the limit.
+ memory_context->set_mem_limit(10);
+ Block block
{ColumnHelper::create_column_with_name<DataTypeInt64>({100, 1000, 1000}),
+ ColumnHelper::create_column_with_name<DataTypeInt64>({2,
3, 3})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(memory_use_limit->value(), 2); // 10 / 5 tasks, floor capped
by the share
+ EXPECT_EQ(local_state->_get_hash_table_size(), 2);
+ EXPECT_EQ(local_state->_pre_aggregated_block->rows(), 3);
+
+ // Pass-through rows are not counted against the pushed-down limit (no
eos), but they
+ // are still reported as processed rows.
+ const int64_t process_rows_before = io_context->process_rows();
+ Block out;
+ bool eos = false;
+ EXPECT_TRUE(op->pull(state.get(), &out, &eos).ok());
+ EXPECT_EQ(out.rows(), 3);
+ EXPECT_FALSE(eos);
+ EXPECT_EQ(io_context->process_rows(), process_rows_before + 3);
+ EXPECT_TRUE(op->need_more_input_data(state.get()));
+ }
+
+ {
+ // Query limit restored: aggregation resumes against the retained hash
table.
+ memory_context->set_mem_limit(5000LL * 1024 * 1024);
+ Block block
{ColumnHelper::create_column_with_name<DataTypeInt64>({1000, 1}),
+ ColumnHelper::create_column_with_name<DataTypeInt64>({3,
4})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 4);
+ EXPECT_EQ(local_state->_pre_aggregated_block->rows(), 0);
+ }
+
+ { EXPECT_TRUE(local_state->close(state.get()).ok()); }
+}
+
+TEST_F(StreamingAggOperatorTest,
sort_limit_boundary_survives_pass_through_and_recover) {
+ // TopN(k ASC, 2) pushed down as a sort limit on the group key, with a
real aggregate so
+ // that the pass-through serialization path is exercised. The mock sum
reads column 0 and
+ // the group key slot reads column 1.
+ op->_aggregate_evaluators.push_back(create_agg_fn(pool, "sum",
+
{std::make_shared<DataTypeInt64>()},
+
std::make_shared<DataTypeInt64>(), false));
+ op->_pool = &pool;
+ op->_needs_finalize = false;
+ op->_do_sort_limit = true;
+ op->_sort_limit = 2;
+ op->_order_directions = {1};
+ op->_null_directions = {1};
+ op->set_parallel_tasks(5);
+
+ EXPECT_TRUE(op->set_child(child_op));
+ EXPECT_TRUE(op->prepare(state.get()).ok());
+ op->_probe_expr_ctxs = MockSlotRef::create_mock_contexts(1,
std::make_shared<DataTypeInt64>());
+
+ {
+ auto local_state =
std::make_unique<MockStreamingAggLocalState>(state.get(), op.get());
+ LocalStateInfo info {.parent_profile = &profile,
+ .scan_ranges = {},
+ .shared_state = nullptr,
+ .shared_state_map = {},
+ .task_idx = 0};
+ EXPECT_TRUE(local_state->init(state.get(), info).ok());
+ state->resize_op_id_to_local_state(-100);
+ state->emplace_local_state(op->operator_id(), std::move(local_state));
+ }
+ local_state =
+
static_cast<MockStreamingAggLocalState*>(state->get_local_state(op->operator_id()));
+ EXPECT_TRUE(local_state->open(state.get()).ok());
+ local_state->use_real_decision = true;
+ auto* memory_context =
state->get_query_ctx()->resource_ctx()->memory_context();
+
+ auto boundary_key = [&]() {
+ return
local_state->limit_columns[0]->get_int(local_state->limit_columns_min);
+ };
+
+ {
+ // Two groups seed the heap: candidates {10, 20}, boundary 20.
+ memory_context->set_mem_limit(5000LL * 1024 * 1024);
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1,
1}),
+ ColumnHelper::create_column_with_name<DataTypeInt64>({10,
20})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 2);
+ EXPECT_EQ(local_state->need_do_sort_limit, 1);
+ EXPECT_EQ(boundary_key(), 20);
+ }
+
+ {
+ // Query limit lowered: the mixed block passes through. Key 30 is
beyond the boundary
+ // and is dropped; key 5 is forwarded without entering the hash table,
so it must not
+ // tighten the boundary either.
+ memory_context->set_mem_limit(10);
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1,
1}),
+ ColumnHelper::create_column_with_name<DataTypeInt64>({5,
30})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 2);
+ EXPECT_EQ(local_state->_pre_aggregated_block->rows(), 1);
+
EXPECT_EQ(local_state->_pre_aggregated_block->get_by_position(0).column->get_int(0),
5);
+ EXPECT_EQ(boundary_key(), 20);
+
+ Block out;
+ bool eos = false;
+ EXPECT_TRUE(op->pull(state.get(), &out, &eos).ok());
+ EXPECT_EQ(out.rows(), 1);
+ EXPECT_FALSE(eos);
+ }
+
+ {
+ // Query limit restored: key 5 is new to the hash table and enters the
heap exactly
+ // once, so the candidates become {5, 10} with boundary 10, not {5, 5}.
+ memory_context->set_mem_limit(5000LL * 1024 * 1024);
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1}),
+
ColumnHelper::create_column_with_name<DataTypeInt64>({5})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 3);
+ EXPECT_EQ(local_state->_pre_aggregated_block->rows(), 0);
+ EXPECT_EQ(boundary_key(), 10);
+ }
+
+ {
+ // Key 7 belongs in the TopN result {5, 7}: it must be kept, and it
tightens the
+ // boundary to 7.
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1}),
+
ColumnHelper::create_column_with_name<DataTypeInt64>({7})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 4);
+ EXPECT_EQ(boundary_key(), 7);
+ }
+
+ {
+ // Key 8 is beyond the boundary and is dropped without touching the
hash table.
+ Block block {ColumnHelper::create_column_with_name<DataTypeInt64>({1}),
+
ColumnHelper::create_column_with_name<DataTypeInt64>({8})};
+ auto st = op->push(state.get(), &block, false);
+ EXPECT_TRUE(st.ok()) << st.msg();
+ EXPECT_EQ(local_state->_get_hash_table_size(), 4);
+ EXPECT_EQ(local_state->_pre_aggregated_block->rows(), 0);
+ EXPECT_EQ(boundary_key(), 7);
}
{ EXPECT_TRUE(local_state->close(state.get()).ok()); }
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
index ddaec7b2ae9..6fb6656ce76 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/SessionVariable.java
@@ -3280,9 +3280,11 @@ public class SessionVariable implements Serializable,
Writable {
@VarAttrDef.VarAttr(name = LOW_MEMORY_MODE_BUFFER_LIMIT, fuzzy = false)
public long lowMemoryModeBufferLimit = 33554432;
- // The memory limit of streaming agg when spilling is enabled
+ // The memory limit of streaming agg when spilling is enabled. It is
applied on top of the
+ // query-limit-based budget (1/5 of the query memory limit shared among
the tasks);
+ // 0 disables this explicit bound.
// NOTE: streaming agg operator will not spill to disk.
- @VarAttrDef.VarAttr(name = SPILL_STREAMING_AGG_MEM_LIMIT, fuzzy = false)
+ @VarAttrDef.VarAttr(name = SPILL_STREAMING_AGG_MEM_LIMIT, needForward =
true, fuzzy = false)
public long spillStreamingAggMemLimit = 268435456; //256MB
@VarAttrDef.VarAttr(name = SPILL_HASH_JOIN_PARTITION_COUNT, fuzzy = true)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]