This is an automated email from the ASF dual-hosted git repository.

kakachen pushed a commit to branch time_sharing_executor_master
in repository https://gitbox.apache.org/repos/asf/doris.git

commit 73fa83b4bb58aebd0aeec6db1118842890e57f43
Author: kakachen <[email protected]>
AuthorDate: Tue Jul 15 19:49:41 2025 +0800

    Change split queue to not thread safe, becasue it guarded by lock in 
TimeSharingTaskExecutor.
---
 .../simulator/simulation_fifo_split_queue.h        | 41 +------------
 .../time_sharing/multilevel_split_queue.cpp        | 67 +++++-----------------
 .../executor/time_sharing/multilevel_split_queue.h | 14 +----
 .../vec/exec/executor/time_sharing/split_queue.h   |  1 -
 4 files changed, 20 insertions(+), 103 deletions(-)

diff --git a/be/src/vec/exec/executor/simulator/simulation_fifo_split_queue.h 
b/be/src/vec/exec/executor/simulator/simulation_fifo_split_queue.h
index 846d7d6611b..6e33fada0c1 100644
--- a/be/src/vec/exec/executor/simulator/simulation_fifo_split_queue.h
+++ b/be/src/vec/exec/executor/simulator/simulation_fifo_split_queue.h
@@ -34,34 +34,18 @@ public:
 
     int compute_level(int64_t scheduled_nanos) override { return 0; }
 
-    void offer(std::shared_ptr<PrioritizedSplitRunner> split) override {
-        {
-            std::lock_guard<std::mutex> lock(_mutex);
-            _queue.push(split);
-        }
-        _not_empty.notify_one();
-    }
+    void offer(std::shared_ptr<PrioritizedSplitRunner> split) override { 
_queue.push(split); }
 
     std::shared_ptr<PrioritizedSplitRunner> take() override {
-        std::unique_lock<std::mutex> lock(_mutex);
-        _not_empty.wait(lock, [this] { return !_queue.empty() || _interrupted; 
});
-
-        if (_interrupted) {
-            return nullptr;
-        }
-
+        if (_queue.empty()) return nullptr;
         auto split = _queue.front();
         _queue.pop();
         return split;
     }
 
-    size_t size() const override {
-        std::lock_guard<std::mutex> lock(_mutex);
-        return _queue.size();
-    }
+    size_t size() const override { return _queue.size(); }
 
     void remove(std::shared_ptr<PrioritizedSplitRunner> split) override {
-        std::lock_guard<std::mutex> lock(_mutex);
         std::queue<std::shared_ptr<PrioritizedSplitRunner>> new_queue;
         while (!_queue.empty()) {
             auto current = _queue.front();
@@ -71,13 +55,9 @@ public:
             }
         }
         _queue.swap(new_queue);
-        if (_queue.empty()) {
-            _not_empty.notify_all();
-        }
     }
 
     void remove_all(const 
std::vector<std::shared_ptr<PrioritizedSplitRunner>>& splits) override {
-        std::lock_guard<std::mutex> lock(_mutex);
         std::unordered_set<std::shared_ptr<PrioritizedSplitRunner>> 
to_remove(splits.begin(),
                                                                               
splits.end());
         std::queue<std::shared_ptr<PrioritizedSplitRunner>> new_queue;
@@ -89,26 +69,14 @@ public:
             }
         }
         _queue.swap(new_queue);
-        if (_queue.empty()) {
-            _not_empty.notify_all();
-        }
     }
 
     void clear() override {
-        std::lock_guard<std::mutex> lock(_mutex);
         while (!_queue.empty()) {
             _queue.pop();
         }
     }
 
-    void interrupt() override {
-        {
-            std::lock_guard<std::mutex> lock(_mutex);
-            _interrupted = true;
-        }
-        _not_empty.notify_all();
-    }
-
     Priority update_priority(const Priority& old_priority, int64_t 
quanta_nanos,
                              int64_t scheduled_nanos) override {
         return old_priority;
@@ -120,9 +88,6 @@ public:
 
 private:
     std::queue<std::shared_ptr<PrioritizedSplitRunner>> _queue;
-    mutable std::mutex _mutex;
-    std::condition_variable _not_empty;
-    std::atomic<bool> _interrupted {false};
 };
 
 } // namespace vectorized
diff --git a/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.cpp 
b/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.cpp
index 340dee2d53a..f27e9ba1a62 100644
--- a/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.cpp
+++ b/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.cpp
@@ -99,44 +99,30 @@ int64_t MultilevelSplitQueue::get_level_min_priority(int 
level, int64_t task_thr
 void MultilevelSplitQueue::offer(std::shared_ptr<PrioritizedSplitRunner> 
split) {
     split->set_ready();
     int level = split->priority().level();
-    std::unique_lock<std::mutex> lock(_mutex);
-    _offer_locked(split, level, lock);
+    _do_offer(split, level);
 }
 
-void 
MultilevelSplitQueue::_offer_locked(std::shared_ptr<PrioritizedSplitRunner> 
split, int level,
-                                         std::unique_lock<std::mutex>& lock) {
+void MultilevelSplitQueue::_do_offer(std::shared_ptr<PrioritizedSplitRunner> 
split, int level) {
     if (_level_waiting_splits[level].empty()) {
-        // Accesses to _level_scheduled_time are not synchronized, so we have 
a data race
-        // here - our level time math will be off. However, the staleness is 
bounded by
-        // the fact that only running splits that complete during this 
computation
-        // can update the level time. Therefore, this is benign.
-        int64_t level0_time = _get_level0_target_time(lock);
+        int64_t level0_time = _get_level0_target_time();
         int64_t level_expected_time =
                 static_cast<int64_t>(level0_time / 
std::pow(_level_time_multiplier, level));
         int64_t delta = level_expected_time - 
_level_scheduled_time[level].load();
         _level_scheduled_time[level].fetch_add(delta);
     }
-
     _level_waiting_splits[level].push(split);
-    _not_empty.notify_all();
 }
 
 std::shared_ptr<PrioritizedSplitRunner> MultilevelSplitQueue::take() {
-    std::unique_lock<std::mutex> lock(_mutex);
-
-    while (!_interrupted) {
-        auto split = _poll_split(lock);
-        if (split) {
-            if (split->update_level_priority()) {
-                _offer_locked(split, split->priority().level(), lock);
-                continue;
-            }
-            int selected_level = split->priority().level();
-            
_level_min_priority[selected_level].store(split->priority().level_priority());
-            return split;
+    auto split = _poll_split();
+    if (split) {
+        if (split->update_level_priority()) {
+            _do_offer(split, split->priority().level());
+            return take();
         }
-
-        _not_empty.wait(lock);
+        int selected_level = split->priority().level();
+        
_level_min_priority[selected_level].store(split->priority().level_priority());
+        return split;
     }
     return nullptr;
 }
@@ -149,9 +135,8 @@ std::shared_ptr<PrioritizedSplitRunner> 
MultilevelSplitQueue::take() {
  * with the objective of minimizing deviation from the target scheduled time. 
From this level,
  * we pick the split with the lowest priority.
  */
-std::shared_ptr<PrioritizedSplitRunner> MultilevelSplitQueue::_poll_split(
-        std::unique_lock<std::mutex>& lock) {
-    int64_t target_scheduled_time = _get_level0_target_time(lock);
+std::shared_ptr<PrioritizedSplitRunner> MultilevelSplitQueue::_poll_split() {
+    int64_t target_scheduled_time = _get_level0_target_time();
     double worst_ratio = 1.0;
     int selected_level = -1;
 
@@ -179,14 +164,11 @@ std::shared_ptr<PrioritizedSplitRunner> 
MultilevelSplitQueue::_poll_split(
 }
 
 void MultilevelSplitQueue::remove(std::shared_ptr<PrioritizedSplitRunner> 
split) {
-    std::lock_guard<std::mutex> lock(_mutex);
-
     for (auto& level_queue : _level_waiting_splits) {
         std::priority_queue<std::shared_ptr<PrioritizedSplitRunner>,
                             
std::vector<std::shared_ptr<PrioritizedSplitRunner>>,
                             SplitRunnerComparator>
                 new_queue;
-
         while (!level_queue.empty()) {
             auto current = level_queue.top();
             level_queue.pop();
@@ -196,17 +178,10 @@ void 
MultilevelSplitQueue::remove(std::shared_ptr<PrioritizedSplitRunner> split)
         }
         level_queue.swap(new_queue);
     }
-
-    if (std::all_of(_level_waiting_splits.begin(), _level_waiting_splits.end(),
-                    [](const auto& q) { return q.empty(); })) {
-        _not_empty.notify_all();
-    }
 }
 
 void MultilevelSplitQueue::remove_all(
         const std::vector<std::shared_ptr<PrioritizedSplitRunner>>& splits) {
-    std::lock_guard<std::mutex> lock(_mutex);
-
     std::unordered_set<std::shared_ptr<PrioritizedSplitRunner>> 
to_remove(splits.begin(),
                                                                           
splits.end());
 
@@ -215,7 +190,6 @@ void MultilevelSplitQueue::remove_all(
                             
std::vector<std::shared_ptr<PrioritizedSplitRunner>>,
                             SplitRunnerComparator>
                 new_queue;
-
         while (!level_queue.empty()) {
             auto current = level_queue.top();
             level_queue.pop();
@@ -225,26 +199,14 @@ void MultilevelSplitQueue::remove_all(
         }
         level_queue.swap(new_queue);
     }
-
-    if (std::all_of(_level_waiting_splits.begin(), _level_waiting_splits.end(),
-                    [](const auto& q) { return q.empty(); })) {
-        _not_empty.notify_all();
-    }
 }
 
 size_t MultilevelSplitQueue::size() const {
-    std::lock_guard<std::mutex> lock(_mutex);
     return std::accumulate(_level_waiting_splits.begin(), 
_level_waiting_splits.end(), size_t(0),
                            [](size_t sum, const auto& queue) { return sum + 
queue.size(); });
 }
 
-void MultilevelSplitQueue::interrupt() {
-    std::lock_guard<std::mutex> lock(_mutex);
-    _interrupted = true;
-    _not_empty.notify_all();
-}
-
-int64_t 
MultilevelSplitQueue::_get_level0_target_time(std::unique_lock<std::mutex>& 
lock) {
+int64_t MultilevelSplitQueue::_get_level0_target_time() {
     int64_t level0_target_time = _level_scheduled_time[0].load();
     double current_multiplier = _level_time_multiplier;
 
@@ -258,7 +220,6 @@ int64_t 
MultilevelSplitQueue::_get_level0_target_time(std::unique_lock<std::mute
 }
 
 void MultilevelSplitQueue::clear() {
-    std::lock_guard<std::mutex> lock(_mutex);
     for (auto& queue : _level_waiting_splits) {
         while (!queue.empty()) {
             queue.pop();
diff --git a/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.h 
b/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.h
index c3363a42c79..14d1363f6d7 100644
--- a/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.h
+++ b/be/src/vec/exec/executor/time_sharing/multilevel_split_queue.h
@@ -17,9 +17,7 @@
 
 #pragma once
 #include <array>
-#include <condition_variable>
 #include <memory>
-#include <mutex>
 #include <queue>
 
 #include "common/factory_creator.h"
@@ -79,15 +77,13 @@ public:
     void remove(std::shared_ptr<PrioritizedSplitRunner> split) override;
     void remove_all(const 
std::vector<std::shared_ptr<PrioritizedSplitRunner>>& splits) override;
     void clear() override;
-    void interrupt() override;
 
     int64_t level_scheduled_time(int level) const { return 
_level_scheduled_time[level].load(); }
 
 private:
-    int64_t _get_level0_target_time(std::unique_lock<std::mutex>& lock);
-    std::shared_ptr<PrioritizedSplitRunner> 
_poll_split(std::unique_lock<std::mutex>& lock);
-    void _offer_locked(std::shared_ptr<PrioritizedSplitRunner> split, int 
level,
-                       std::unique_lock<std::mutex>& lock);
+    int64_t _get_level0_target_time();
+    std::shared_ptr<PrioritizedSplitRunner> _poll_split();
+    void _do_offer(std::shared_ptr<PrioritizedSplitRunner> split, int level);
 
     const double _level_time_multiplier;
 
@@ -99,10 +95,6 @@ private:
 
     std::array<std::atomic<int64_t>, LEVEL_THRESHOLD_SECONDS.size()> 
_level_scheduled_time;
     std::array<std::atomic<int64_t>, LEVEL_THRESHOLD_SECONDS.size()> 
_level_min_priority;
-
-    std::atomic<bool> _interrupted {false};
-    mutable std::mutex _mutex;
-    std::condition_variable _not_empty;
 };
 
 } // namespace vectorized
diff --git a/be/src/vec/exec/executor/time_sharing/split_queue.h 
b/be/src/vec/exec/executor/time_sharing/split_queue.h
index 16a96c03a71..b8b6ccd9f2e 100644
--- a/be/src/vec/exec/executor/time_sharing/split_queue.h
+++ b/be/src/vec/exec/executor/time_sharing/split_queue.h
@@ -37,7 +37,6 @@ public:
     virtual void remove(std::shared_ptr<PrioritizedSplitRunner> split) = 0;
     virtual void remove_all(const 
std::vector<std::shared_ptr<PrioritizedSplitRunner>>& splits) = 0;
     virtual void clear() = 0;
-    virtual void interrupt() = 0;
 
     virtual Priority update_priority(const Priority& old_priority, int64_t 
quanta_nanos,
                                      int64_t scheduled_nanos) = 0;


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to