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]

Reply via email to