mrhhsg commented on code in PR #67845:
URL: https://github.com/apache/doris/pull/67845#discussion_r4228617093
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -300,6 +302,19 @@ void WorkloadGroupMgr::add_paused_query(const
std::shared_ptr<ResourceContext>&
LOG(INFO) << "Insert one new paused query: "
<< resource_ctx->task_controller()->debug_string()
<< ", workload group: " << wg->debug_string();
+ } else {
+ // Another task of the same query already failed a reservation. The
query is resumed as
+ // a whole and every blocked task then retries its own request, so
record the sum of the
+ // pending requests: the query is woken up only once all of them fit
at the same time.
+ // Waking it up as soon as one of them fits would fail the others
again and start a new
+ // wait from scratch.
+ LOG(INFO) << "Query: " <<
print_id(resource_ctx->task_controller()->task_id())
+ << " is already paused, add the pending reservation "
+ << PrettyPrinter::print_bytes(reserve_size) << " to the
recorded "
+ << PrettyPrinter::print_bytes(it->reserve_size_);
+ auto node = queries_list.extract(it);
+ node.value().reserve_size_ += reserve_size;
Review Comment:
Fixed in 65ba0b876f7. `add_paused_query()` records the largest pending
reservation again (as in a428b3e787d), so the query is woken up once that one
fits and the tasks reserve one after another (a task releases its reservation
after each block, so the pending requests do not have to fit at the same time).
The wait is now bounded across retries instead of per entry:
`TaskController::start_process_memory_wait()` keeps the time of the first
reservation that failed since the last successful one, `add_paused_query()`
uses it as the `enqueue_at` of a `PROCESS_MEMORY_EXCEEDED` entry, and
`PipelineTask::_try_to_reserve_memory()` ends the wait when a reservation
succeeds. A query resumed to retry that fails again before any of its
reservations succeeded therefore continues the same wait and is still cancelled
at `spill_in_paused_queue_timeout_ms`; a smaller sibling that does not fit next
to the larger request can no longer restart it. `PausedQuery::enqueue_at` is a
`MonotonicMillis()` time
stamp now so that both use the same clock. Tests:
`process_mem_exceeded_keeps_largest_pending_reservation` (64 MiB + 4 MiB
pending, 65 MiB headroom resumes the query although 68 MiB does not fit),
`process_mem_exceeded_largest_pending_reservation_cancels_after_timeout`,
`process_mem_exceeded_failed_retry_continues_the_wait` (resumed past the
timeout, re-paused by the failed retry, cancelled in the next round; without
the change the new entry starts at 0 and stays paused),
`process_mem_exceeded_successful_reservation_starts_new_wait`, and
`PipelineTaskTest.TEST_RESERVE_MEMORY_FAIL_SPILLABLE` asserts that a successful
reservation clears the wait start (fails without the change).
##########
be/test/testutil/mock/mock_query_task_controller.h:
##########
@@ -35,6 +38,52 @@ struct MockQueryTaskController : public QueryTaskController {
}
void set_cancelled_time(int64_t ctime) { cancelled_time_ = ctime; }
+
+ // Memory reclamation asks every candidate query whether it is cancelled
while it scans a
+ // workload group, so this is where a test can let another query release
memory during that
+ // scan.
+ bool is_cancelled() const override {
+ if (on_is_cancelled_) {
+ on_is_cancelled_();
+ }
+ return QueryTaskController::is_cancelled();
+ }
+
+ // Pipeline state seen by WorkloadGroupMgr::handle_single_query_. The
query context of a
+ // unit test has no fragments, so the real implementation never reports a
running or a
+ // revocable task; these knobs simulate them.
+ // NOLINTNEXTLINE(readability-make-member-function-const): overrides a
non-const virtual.
+ void get_revocable_info(size_t* revocable_size, size_t* memory_usage,
+ bool* has_running_task) override {
+ QueryTaskController::get_revocable_info(revocable_size, memory_usage,
has_running_task);
+ *has_running_task = has_running_task_;
+ }
+
+ // The manager only checks whether the list is empty before it calls
revoke_memory(), which
+ // is mocked below, so a placeholder entry is enough to stand for a
revocable task.
+ std::vector<PipelineTask*> get_revocable_tasks() override {
+ if (on_get_revocable_tasks_) {
+ on_get_revocable_tasks_();
+ }
+ return has_revocable_task_ ? std::vector<PipelineTask*> {nullptr}
+ : std::vector<PipelineTask*> {};
+ }
+
+ Status revoke_memory() override {
Review Comment:
Fixed in 65ba0b876f7. `MockQueryTaskController::revoke_memory()` now
completes the spill the way the real implementation does: it clears
`has_revocable_task_` (the spilled tasks have nothing left to revoke) and calls
`set_memory_sufficient(true)`, which readies the memory dependency and resets
the paused reason. `process_mem_exceeded_spills_revocable_tasks_at_hard_limit`
asserts that the dependency is blocked after the pause, that the first round
calls `revoke_memory()` once, resumes the query (`paused_reason().ok()`,
dependency ready) and removes the entry, then pauses it again as the failed
retry (dependency blocked again, no revocable task left) and checks that the
next round cancels it at the hard limit without another `revoke_memory()` call.
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -678,32 +702,50 @@ int64_t
WorkloadGroupMgr::revoke_memory_from_other_groups_() {
// then not revoke memory from it.
continue;
}
- if (total_used_memory - min_memory_limit > max_exceeded_memory) {
- max_wg = workload_group.second;
- max_exceeded_memory = total_used_memory - min_memory_limit;
- }
+ exceeded_wgs.emplace_back(total_used_memory - min_memory_limit,
workload_group.second);
}
}
- if (max_wg == nullptr) {
- return 0;
- }
- if (max_exceeded_memory < 1 << 27) {
- LOG(INFO) << "The workload group that exceed most memory is :"
- << max_wg->memory_debug_string() << ", max_exceeded_memory: "
- << PrettyPrinter::print(max_exceeded_memory, TUnit::BYTES)
- << " less than 128MB, no need to revoke memory";
- return 0;
- }
- int64_t freed_mem = static_cast<int64_t>((double)max_exceeded_memory *
0.1);
- // Revoke 10% of memory from the workload group that exceed most memory
- max_wg->revoke_memory(freed_mem, "exceed_memory", profile.get());
- std::stringstream ss;
- profile->pretty_print(&ss);
- LOG(INFO) << fmt::format(
- "[MemoryGC] process memory not enough, revoke memory from
workload_group: {}, "
- "free memory {}. cost(us): {}, details: {}",
- max_wg->memory_debug_string(),
PrettyPrinter::print_bytes(freed_mem),
- watch.elapsed_time() / 1000, ss.str());
+ std::sort(exceeded_wgs.begin(), exceeded_wgs.end(),
+ [](const auto& lhs, const auto& rhs) { return lhs.first >
rhs.first; });
+
+ int64_t freed_mem = 0;
+ size_t tried_wgs = 0;
+ for (const auto& [exceeded_memory, wg] : exceeded_wgs) {
+ if (exceeded_memory < 1 << 27) {
+ // The remaining workload groups exceed even less.
+ LOG(INFO) << "The workload group that exceed most memory among the
untried ones is :"
+ << wg->memory_debug_string() << ", exceeded_memory: "
+ << PrettyPrinter::print(exceeded_memory, TUnit::BYTES)
+ << " less than 128MB, no need to revoke memory";
+ break;
+ }
Review Comment:
On the variant carried in the latest review (a peer falling below its
minimum between the walk's `refresh_memory_usage()` and
`WorkloadGroup::revoke_memory()`): the window left is the few calls between
that refresh and `memory_used()` at the start of `revoke_memory()`, with no
waiting in between; the window that mattered (scanning the earlier peers for up
to `revoke_memory_max_tolerance_ms`) is what the refresh in the walk closes.
The same check-then-act window exists inside
`MemoryReclamation::revoke_tasks_memory()` itself, which snapshots the
consumption and then cancels, and `WorkloadGroup::revoke_memory()` is shared
with the memory gc callers; adding a min-memory eligibility check inside it
would change their semantics, so I am leaving this as is in this PR.
##########
be/src/runtime/workload_group/workload_group_manager.cpp:
##########
@@ -824,38 +839,45 @@ bool WorkloadGroupMgr::handle_single_query_(const
std::shared_ptr<ResourceContex
return true;
}
} else {
- // Should not consider about process memory. For example, the query's
limit is 100g, workload
- // group's memlimit is 10g, process memory is 20g. The query reserve
will always failed in wg
- // limit, and process is always have memory, so that it will resume
and failed reserve again.
- const size_t test_memory_size = std::max<size_t>(size_to_reserve, 32L
* 1024 * 1024);
- if
(!GlobalMemoryArbitrator::is_exceed_soft_mem_limit(test_memory_size)) {
- LOG(INFO) << "Query: " << query_id
- << ", process limit not exceeded now, resume this query"
- << ", process memory info: "
- <<
GlobalMemoryArbitrator::process_memory_used_details_str()
- << ", wg info: " << wg->debug_string();
- requestor->task_controller()->set_memory_sufficient(true);
- return true;
- } else {
- // if cannot find any memory to release, then let the query
continue to run as far as possible
- // or cancelled by gc if memory is really not enough.
- Status error_status = Status::MemoryLimitExceeded(
- "Query {} process memory is exceeded"
- ", and there is no cache now. And could not find task to
spill, disable "
- "reserve memory and resume it. "
- "Query memory usage: {}, limit: {}, reserved "
- "size: {}, try to reserve: {}, wg info: {}."
- " Maybe you should set the workload group's limit to a
lower value. {}",
- query_id, PrettyPrinter::print_bytes(memory_usage),
- PrettyPrinter::print_bytes(limit),
PrettyPrinter::print_bytes(reserved_size),
- PrettyPrinter::print_bytes(size_to_reserve),
wg->memory_debug_string(),
- doris::ProcessProfile::instance()
- ->memory_profile()
- ->process_memory_detail_str());
- LOG_LONG_STRING(INFO, error_status.to_string());
- requestor->task_controller()->cancel(error_status);
- return true;
+ // PROCESS_MEMORY_EXCEEDED. The caller
(handle_process_memory_exceeded_) has already
+ // resumed the query if the process is no longer above the soft memory
limit, so the
Review Comment:
On the QUERY-first variant carried in the latest review:
`TaskController::update_paused_reason()` keeps `QUERY_MEMORY_EXCEEDED` over a
later `PROCESS_MEMORY_EXCEEDED` failure of a sibling task, which is
pre-existing precedence that this PR does not change. A QUERY reason means the
query is over its own limit, and the QUERY route either spills it or resumes it
with `disable_reserve_memory()`; from then on every allocation goes through
`Allocator::sys_memory_check()`, which checks `is_exceed_hard_mem_limit()`
directly and fails the allocation regardless of `disable_memory_gc`. The
running-task exit on that route lasts for the current operator call only. The
mixed-reason handling belongs to the existing QUERY route rather than to the
PROCESS route this PR changes, so no change here.
--
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]