github-actions[bot] commented on code in PR #68713:
URL: https://github.com/apache/doris/pull/68713#discussion_r4234234109
##########
be/src/exec/scan/scanner_context.cpp:
##########
@@ -362,6 +363,61 @@ void
ScannerContext::push_completed_scan_task(std::shared_ptr<ScanTask> scan_tas
_dependency->set_ready();
}
+void ScannerContext::park_scan_task(std::shared_ptr<ScanTask> scan_task,
+ SharedListenableFuture<Void> waiting_for) {
+ {
+ std::unique_lock<std::mutex> l(_transfer_lock);
+ scan_task->set_state(ScanTask::State::PARKED);
+ _in_flight_tasks_num--;
+ ++_parked_tasks_num;
+ // The slot it held may now admit a pending scanner - one whose block
the operator consumed
+ // while this context had no slot to spare, holding what this task
waits for, among them.
+ // Nothing else would look: a parked task never completes a scan
attempt to make the operator
+ // schedule again.
+ if (!done()) {
+ Status status =
_scanner_scheduler->schedule_scan_task(shared_from_this(), nullptr, l);
+ if (!status.ok()) {
+ set_context_failure(status, l);
+ }
+ }
+ }
+ std::weak_ptr<ScannerContext> weak_ctx = shared_from_this();
+ // Runs on the thread that completes the future, or right here if it is
done already.
+ waiting_for.add_callback([weak_ctx, scan_task](const Void&, const Status&)
{
+ if (auto ctx = weak_ctx.lock()) {
+ ctx->_resume_parked_task(scan_task);
+ }
+ });
+}
+
+void ScannerContext::_resume_parked_task(const std::shared_ptr<ScanTask>&
scan_task) {
+ auto task_execution_lock = task_exec_ctx();
+ if (task_execution_lock == nullptr) {
+ // The query has finished; nothing will read this task again.
+ return;
+ }
+#ifndef BE_TEST
+ // Scheduling allocates for the query, not for the gate's thread that
completed the future. A
+ // future done before park_scan_task() added this callback runs it on the
worker that parked
+ // the task, which is attached to this query already and must stay so.
+ std::optional<AttachTask> attach;
+ if (!thread_context()->is_attach_task()) {
+ attach.emplace(_state);
+ }
+#endif
+ std::unique_lock<std::mutex> l(_transfer_lock);
+ --_parked_tasks_num;
+ if (done()) {
+ return;
+ }
+ // Like a task whose block the operator just consumed: it may run again.
+ scan_task->set_state(ScanTask::State::PENDING);
+ Status status = _scanner_scheduler->schedule_scan_task(shared_from_this(),
scan_task, l);
Review Comment:
[P1] Finish the current scan invocation before resubmitting this parked
task. The gate can complete its future before `park_scan_task()` registers the
callback; `add_callback()` then runs it inline, and this call can schedule the
same `ScanTask` on another worker before the first `_scanner_scan()` returns.
`execute_scan_task()` still reads the task's plain `is_eos()` state after that
return while the new worker changes it. On TaskExecutor this can also requeue
the same `ScannerSplitRunner` before its first `process_for()` exits, allowing
concurrent execution and premature split cleanup. Defer resubmission until the
original invocation has fully exited, including the already-ready future case.
##########
be/src/util/jni_scan_heap_gate.cpp:
##########
@@ -0,0 +1,256 @@
+// 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.
+
+#include "util/jni_scan_heap_gate.h"
+
+#include <algorithm>
+#include <chrono>
+#include <limits>
+#include <utility>
+
+#include "common/config.h"
+#include "common/logging.h"
+#include "util/jni-util.h"
+#include "util/thread.h"
+#include "util/time.h"
+
+namespace doris {
+
+namespace {
+
+// How often the gate asks the waiting readers whether their scans stopped,
and looks for those that
+// waited too long, although nobody gave a share back.
+constexpr int64_t POLL_INTERVAL_NS = 100L * 1000 * 1000;
+constexpr int64_t MB = 1024 * 1024;
+
+// A share of the -Xmx the JVM was started with, read from the same options
the JVM was created from
+// - not measured, so nothing here calls into the JVM.
+int64_t jvm_heap_budget() {
+ const double budget =
static_cast<double>(Jni::Util::get_max_jni_heap_memory_size()) *
+ config::jni_scanner_heap_budget_ratio;
Review Comment:
[P2] Synchronize reads of both mutable heap-gate settings. `DEFINE_mDouble`
and `DEFINE_mInt64` make `jni_scanner_heap_budget_ratio` and
`jni_scanner_heap_max_wait_ms` ordinary fields. `config::set_config()` writes
them under `mutable_string_config_lock`, but the gate reads the ratio here and
the timeout in `_poll()` without that lock while its thread handles waiters.
Updating either setting online therefore races admission or timeout decisions.
Read synchronized snapshots (or use atomic storage) for both values.
--
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]