diff --git a/src/spider/CMakeLists.txt b/src/spider/CMakeLists.txt index a17adcaa7..61e735a08 100644 --- a/src/spider/CMakeLists.txt +++ b/src/spider/CMakeLists.txt @@ -115,8 +115,6 @@ set(SPIDER_SCHEDULER_SOURCES scheduler/SchedulerMessage.hpp scheduler/SchedulerServer.cpp scheduler/SchedulerServer.hpp - scheduler/SchedulerTaskCache.cpp - scheduler/SchedulerTaskCache.hpp utils/StopToken.hpp CACHE INTERNAL "spider scheduler source files" diff --git a/src/spider/scheduler/FifoPolicy.cpp b/src/spider/scheduler/FifoPolicy.cpp index bc5fd0382..afd7311ab 100644 --- a/src/spider/scheduler/FifoPolicy.cpp +++ b/src/spider/scheduler/FifoPolicy.cpp @@ -2,12 +2,15 @@ #include #include +#include #include #include #include #include +#include #include +#include #include #include #include @@ -18,16 +21,11 @@ #include "../storage/DataStorage.hpp" #include "../storage/MetadataStorage.hpp" #include "../storage/StorageConnection.hpp" -#include "SchedulerTaskCache.hpp" -namespace { +namespace spider::scheduler { -auto task_locality_satisfied( - std::shared_ptr const& data_store, - spider::core::StorageConnection& conn, - spider::core::Task const& task, - std::string const& addr -) -> bool { +auto FifoPolicy::task_locality_satisfied(spider::core::Task const& task, std::string const& addr) + -> bool { for (auto const& input : task.get_inputs()) { if (input.get_value().has_value()) { continue; @@ -37,11 +35,16 @@ auto task_locality_satisfied( continue; } boost::uuids::uuid const data_id = optional_data_id.value(); - spider::core::Data data; - if (false == data_store->get_data(conn, data_id, &data).success()) { - throw std::runtime_error( - fmt::format("Data with id {} not exists.", boost::uuids::to_string((data_id))) - ); + core::Data data; + if (m_data_cache.contains(data_id)) { + data = m_data_cache[data_id]; + } else { + if (false == m_data_store->get_data(m_conn, data_id, &data).success()) { + throw std::runtime_error( + fmt::format("Data with id {} not exists.", to_string((data_id))) + ); + } + m_data_cache.emplace(data_id, data); } if (false == data.is_hard_locality()) { continue; @@ -57,10 +60,6 @@ auto task_locality_satisfied( return true; } -} // namespace - -namespace spider::scheduler { - FifoPolicy::FifoPolicy( std::shared_ptr const& metadata_store, std::shared_ptr const& data_store, @@ -68,78 +67,70 @@ FifoPolicy::FifoPolicy( ) : m_metadata_store{metadata_store}, m_data_store{data_store}, - m_conn{conn}, - m_task_cache{ - metadata_store, - data_store, - conn, - [&](std::vector& tasks, - boost::uuids::uuid const& worker_id, - std::string const& worker_addr) -> std::optional { - return get_next_task(tasks, worker_id, worker_addr); - } - } {} + m_conn{conn} {} -auto FifoPolicy::get_next_task( - std::vector& tasks, - boost::uuids::uuid const& /*worker_id*/, +auto FifoPolicy::schedule_next( + boost::uuids::uuid const /*worker_id*/, std::string const& worker_addr ) -> std::optional { - std::erase_if(tasks, [this, worker_addr](core::Task const& task) -> bool { - return !task_locality_satisfied(m_data_store, m_conn, task, worker_addr); + if (m_tasks.empty()) { + fetch_tasks(); + if (m_tasks.empty()) { + return std::nullopt; + } + } + auto const reverse_begin = std::reverse_iterator(m_tasks.end()); + auto const reverse_end = std::reverse_iterator(m_tasks.begin()); + auto const it = std::find_if(reverse_begin, reverse_end, [&](core::Task const& task) { + return task_locality_satisfied(task, worker_addr); }); - - if (tasks.empty()) { + if (it == reverse_end) { return std::nullopt; } - - auto const earliest_task = std::ranges::min_element( - tasks, - {}, - [this](core::Task const& task) -> std::chrono::system_clock::time_point { - boost::uuids::uuid const task_id = task.get_id(); - boost::uuids::uuid job_id; - std::optional const optional_job_id - = m_task_job_cache.get(task_id); - if (optional_job_id.has_value()) { - job_id = optional_job_id.value(); - } else { - if (false - == m_metadata_store->get_task_job_id(m_conn, task_id, &job_id).success()) { - throw std::runtime_error(fmt::format( - "Task with id {} not exists.", - boost::uuids::to_string(task_id) - )); - } - m_task_job_cache.put(task_id, job_id); - } - - std::optional const optional_time - = m_job_time_cache.get(job_id); - if (optional_time.has_value()) { - return optional_time.value(); - } - - core::JobMetadata job_metadata; - if (false - == m_metadata_store->get_job_metadata(m_conn, job_id, &job_metadata).success()) - { - throw std::runtime_error(fmt::format( - "Job with id {} not exists.", - boost::uuids::to_string(job_id) - )); - } - m_job_time_cache.put(job_id, job_metadata.get_creation_time()); - return job_metadata.get_creation_time(); - } - ); - - return earliest_task->get_id(); + boost::uuids::uuid const task_id = it->get_id(); + for (core::TaskInput const& input : it->get_inputs()) { + std::optional const data_id = input.get_data_id(); + if (data_id.has_value()) { + m_data_cache.erase(data_id.value()); + } + } + m_tasks.erase(std::next(it).base()); + return task_id; } -auto FifoPolicy::schedule_next(boost::uuids::uuid const worker_id, std::string const& worker_addr) - -> std::optional { - return m_task_cache.get_ready_task(worker_id, worker_addr); +auto FifoPolicy::fetch_tasks() -> void { + m_data_cache.clear(); + m_metadata_store->get_ready_tasks(m_conn, &m_tasks); + std::vector> instances; + m_metadata_store->get_task_timeout(m_conn, &instances); + for (auto const& [instance, task] : instances) { + m_tasks.emplace_back(task); + } + + // Sort tasks based on job creation time in descending order. + // NOLINTNEXTLINE(misc-include-cleaner) + absl::flat_hash_map> + job_metadata_map; + auto get_task_job_creation_time + = [&](boost::uuids::uuid const task_id) -> std::chrono::system_clock::time_point { + boost::uuids::uuid job_id; + if (false == m_metadata_store->get_task_job_id(m_conn, task_id, &job_id).success()) { + throw std::runtime_error(fmt::format("Task with id {} not exists.", to_string(task_id)) + ); + } + if (job_metadata_map.contains(job_id)) { + return job_metadata_map[job_id].get_creation_time(); + } + core::JobMetadata job_metadata; + if (false == m_metadata_store->get_job_metadata(m_conn, job_id, &job_metadata).success()) { + throw std::runtime_error(fmt::format("Job with id {} not exists.", to_string(job_id))); + } + job_metadata_map[job_id] = job_metadata; + return job_metadata.get_creation_time(); + }; + std::ranges::sort(m_tasks, [&](core::Task const& a, core::Task const& b) { + return get_task_job_creation_time(a.get_id()) > get_task_job_creation_time(b.get_id()); + }); } } // namespace spider::scheduler diff --git a/src/spider/scheduler/FifoPolicy.hpp b/src/spider/scheduler/FifoPolicy.hpp index d114875da..aab32c244 100644 --- a/src/spider/scheduler/FifoPolicy.hpp +++ b/src/spider/scheduler/FifoPolicy.hpp @@ -1,21 +1,19 @@ #ifndef SPIDER_SCHEDULER_FIFOPOLICY_HPP #define SPIDER_SCHEDULER_FIFOPOLICY_HPP -#include #include #include #include #include +#include #include #include "../core/Task.hpp" #include "../storage/DataStorage.hpp" #include "../storage/MetadataStorage.hpp" #include "../storage/StorageConnection.hpp" -#include "../utils/LruCache.hpp" #include "SchedulerPolicy.hpp" -#include "SchedulerTaskCache.hpp" namespace spider::scheduler { @@ -25,28 +23,23 @@ class FifoPolicy final : public SchedulerPolicy { std::shared_ptr const& metadata_store, std::shared_ptr const& data_store, core::StorageConnection& conn - ); auto schedule_next(boost::uuids::uuid worker_id, std::string const& worker_addr) -> std::optional override; private: - auto get_next_task( - std::vector& tasks, - boost::uuids::uuid const& worker_id, - std::string const& worker_addr - ) -> std::optional; + auto fetch_tasks() -> void; + auto task_locality_satisfied(core::Task const& task, std::string const& addr) -> bool; std::shared_ptr m_metadata_store; std::shared_ptr m_data_store; // NOLINTNEXTLINE(cppcoreguidelines-avoid-const-or-ref-data-members) core::StorageConnection& m_conn; - SchedulerTaskCache m_task_cache; - - core::LruCache m_task_job_cache; - core::LruCache m_job_time_cache; + std::vector m_tasks; + // NOLINTNEXTLINE(misc-include-cleaner) + absl::flat_hash_map> m_data_cache; }; } // namespace spider::scheduler diff --git a/src/spider/scheduler/SchedulerTaskCache.cpp b/src/spider/scheduler/SchedulerTaskCache.cpp deleted file mode 100644 index 7ec559855..000000000 --- a/src/spider/scheduler/SchedulerTaskCache.cpp +++ /dev/null @@ -1,99 +0,0 @@ -#include "SchedulerTaskCache.hpp" - -#include -#include -#include -#include -#include -#include -#include - -#include - -#include "../core/Task.hpp" - -namespace spider::scheduler { - -namespace { -constexpr int cUpdateCount = 100; -constexpr int cUpdateInterval = 5; // 5 ms -} // namespace - -auto SchedulerTaskCache::get_ready_task( - boost::uuids::uuid const& worker_id, - std::string const& worker_addr -) -> std::optional { - bool const updated = should_fetch_tasks(); - if (updated) { - fetch_ready_tasks(); - } - - std::optional task = pop_next_task(worker_id, worker_addr); - if (task.has_value()) { - m_update_count++; - return task.value().get_id(); - } - if (updated) { - m_update_count++; - return std::nullopt; - } - fetch_ready_tasks(); - task = pop_next_task(worker_id, worker_addr); - m_update_count++; - if (task.has_value()) { - return task.value().get_id(); - } - - return std::nullopt; -} - -auto SchedulerTaskCache::pop_next_task( - boost::uuids::uuid const& worker_id, - std::string const& worker_addr -) -> std::optional { - std::vector tasks; - tasks.reserve(m_tasks.size()); - for (auto const& task : m_tasks) { - tasks.push_back(task.second); - } - std::optional const task_id = std::invoke( - m_get_next_task_function, - std::ref(tasks), - std::cref(worker_id), - std::cref(worker_addr) - ); - if (!task_id.has_value()) { - return std::nullopt; - } - - core::Task task = m_tasks.at(task_id.value()); - m_tasks.erase(task_id.value()); - return task; -} - -auto SchedulerTaskCache::should_fetch_tasks() -> bool { - std::chrono::steady_clock::time_point const now = std::chrono::steady_clock::now(); - if (m_last_update + std::chrono::milliseconds(cUpdateInterval) < now) { - return true; - } - return m_update_count > cUpdateCount; -} - -void SchedulerTaskCache::fetch_ready_tasks() { - m_tasks.clear(); - std::vector tasks; - m_metadata_store->get_ready_tasks(m_conn, &tasks); - for (core::Task const& task : tasks) { - m_tasks.emplace(std::make_pair(task.get_id(), task)); - } - std::vector> task_instances; - m_metadata_store->get_task_timeout(m_conn, &task_instances); - for (auto const& [task_instance, task] : task_instances) { - m_tasks.emplace(std::make_pair(task.get_id(), task)); - } - - m_update_count = 0; - m_last_update = std::chrono::steady_clock::now(); -} - -} // namespace spider::scheduler diff --git a/src/spider/scheduler/SchedulerTaskCache.hpp b/src/spider/scheduler/SchedulerTaskCache.hpp deleted file mode 100644 index dae3d1ffd..000000000 --- a/src/spider/scheduler/SchedulerTaskCache.hpp +++ /dev/null @@ -1,70 +0,0 @@ -#ifndef SPIDER_SCHEDULER_SCHEDULERTASKCACHE_HPP -#define SPIDER_SCHEDULER_SCHEDULERTASKCACHE_HPP - -#include -#include -#include -#include -#include -#include -#include - -#include -#include - -#include "../core/Task.hpp" -#include "../storage/DataStorage.hpp" -#include "../storage/MetadataStorage.hpp" -#include "../storage/StorageConnection.hpp" - -namespace spider::scheduler { - -class SchedulerTaskCache { -public: - SchedulerTaskCache( - std::shared_ptr const& metadata_store, - std::shared_ptr const& data_store, - core::StorageConnection& conn, - std::function( - std::vector& tasks, - boost::uuids::uuid const& worker_id, - std::string const& worker_addr - )> const& get_next_task_function - ) - : m_metadata_store{metadata_store}, - m_data_store{data_store}, - m_conn{conn}, - m_get_next_task_function{get_next_task_function} {} - - auto get_ready_task(boost::uuids::uuid const& worker_id, std::string const& worker_addr) - -> std::optional; - -private: - auto should_fetch_tasks() -> bool; - - void fetch_ready_tasks(); - - auto pop_next_task(boost::uuids::uuid const& worker_id, std::string const& worker_addr) - -> std::optional; - - std::shared_ptr m_metadata_store; - std::shared_ptr m_data_store; - // NOLINTNEXTLINE(cppcoreguidelines-avoid-const-or-ref-data-members) - core::StorageConnection& m_conn; - - // NOLINTNEXTLINE(misc-include-cleaner) - absl::flat_hash_map> m_tasks; - std::chrono::steady_clock::time_point m_last_update; - size_t m_update_count = 0; - - std::function( - std::vector& tasks, - boost::uuids::uuid const& worker_id, - std::string const& worker_addr - )> - m_get_next_task_function; -}; - -} // namespace spider::scheduler - -#endif // SPIDER_SCHEDULER_SCHEDULERTASKCACHE_HPP