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]
