基于Boost.Asio实现带写优先级的读写锁及任务调度问题
我的应用处理三类任务,来源包括异步外部触发(如用户按键)或内部触发(如超时、任务联动、网络事件等)。每个任务都存在长时间等待阶段(如发送网络请求或等待响应),此时可协作挂起并切换至其他任务。所有任务(无论类型)均使用新连接完成单次收发操作,代码示例如下:
http::response<http::dynamic_body> resp; boost::beast::http::request<boost::beast::http::string_body> req; boost::beast::flat_buffer buf; http::async_write(stream, req, yield); http::async_read(stream, buf, resp, yield);
我希望通过同一组线程运行同一个基础io_context来执行所有任务,相关代码如下:
boost::thread_group threads; num_threads = 0; for(; num_threads < kMaxThreads ; num_threads++) { threads.create_thread([] { try { io_context_.run(); } catch (const std::exception& e) { print("exception found"); } }); }
任务说明
- Type1:发送含预获取OAuth2访问令牌的常规HTTPS请求,可并行执行,通过
spawn/io_context调用,代码示例:
boost::asio::io_context io_context_; boost::asio::spawn(io_context_, [](const boost::asio::yield_context &yield) { ... });
- Type2:基于OAuth2的令牌刷新任务,刷新后旧令牌失效。该任务必须串行执行,原因是执行期间无法确认旧令牌是否已在服务器端失效,且可能需要切换服务器,需等待获取新服务器的令牌后再执行其他任务。
- Type3:无需令牌的HTTP请求,可与前两类任务并行执行。
理想情况下Type2任务通过post/strand实现(无异步yield操作),代码示例:
boost::asio::strand<boost::asio::io_context::executor_type> strand_; // 类构造函数中初始化 strand_(io_context_.get_executor()); boost::asio::post(strand_, []() {...});
但因网络API使用协程,且需允许Type3任务运行,需使用spawn/strand,同时要避免Type2任务进入队列,现询问:如何正确实现该逻辑?
另外,当前实现可保证Type1任务并行、Type2任务串行且其他任务等待其完成,现在希望给Type2任务赋予优先级,使其进入队列后在当前执行任务结束后优先执行,请问:如何通过Boost.Asio实现该优先级调度?
一、Type2任务的spawn/strand正确实现
要在协程场景下保证Type2任务串行,同时不阻塞Type1/Type3任务,只需将spawn的目标执行器指定为Type2专用的strand即可:
- 在类中初始化专用strand:
class TaskManager { private: boost::asio::io_context io_context_; boost::asio::strand<boost::asio::io_context::executor_type> token_refresh_strand_; public: TaskManager() : token_refresh_strand_(io_context_.get_executor()) {} // ...其他成员 };
- 提交Type2任务时,直接将strand作为
spawn的第一个参数:
// 提交令牌刷新协程任务 boost::asio::spawn(token_refresh_strand_, [this](boost::asio::yield_context yield) { // 执行令牌刷新的网络逻辑,包括async_write/async_read等协程调用 boost::beast::tcp_stream stream(io_context_); // ...建立连接、构造请求等操作 http::async_write(stream, refresh_req, yield); http::async_read(stream, buf, refresh_resp, yield); // 更新全局令牌、服务器地址等逻辑 update_oauth_token(refresh_resp); });
这样所有Type2任务都会被strand串行化执行,Type1/Type3任务直接提交到io_context,可以和Type2任务并行运行——strand仅限制Type2自身的串行,不会阻塞其他任务的调度。
如果需要确保Type1任务在Type2执行期间暂停(避免使用旧令牌发起请求),可额外引入全局"令牌刷新中"标志:
- Type2任务执行前设置标志为
true - Type1任务发起请求前检查标志,若为
true则提交到strand等待,或挂起等待标志变更 - Type2任务执行完成后重置标志为
false,并唤醒等待的Type1任务
二、Boost.Asio实现Type2任务的优先级调度
Boost.Asio原生io_context的任务队列是FIFO的,要实现优先级调度,可通过以下两种方案:
方案1:自定义优先级任务队列包装io_context
重载io_context::service替换默认调度逻辑,实现优先级任务队列:
#include <boost/asio/io_context.hpp> #include <queue> #include <tuple> #include <mutex> #include <functional> // 定义优先级类型,数字越大优先级越高 using TaskPriority = int; class PriorityIoService : public boost::asio::io_context::service { public: static boost::asio::io_context::id id; explicit PriorityIoService(boost::asio::io_context& io_context) : boost::asio::io_context::service(io_context) {} // 提交带优先级的任务 void post(TaskPriority priority, std::function<void()> task) { std::lock_guard<std::mutex> lock(mutex_); tasks_.emplace(priority, std::move(task)); // 通知io_context有任务待处理 this->get_io_context().post([this] { dispatch_next(); }); } private: void dispatch_next() { std::lock_guard<std::mutex> lock(mutex_); if (!tasks_.empty()) { auto [priority, task] = std::move(tasks_.top()); tasks_.pop(); task(); } } // 优先级队列(大顶堆,高优先级任务先出) std::priority_queue<std::pair<TaskPriority, std::function<void()>>> tasks_; std::mutex mutex_; // service接口必需的虚函数 void shutdown() override { std::lock_guard<std::mutex> lock(mutex_); tasks_ = {}; } }; boost::asio::io_context::id PriorityIoService::id;
使用方式:
// 初始化优先级服务 PriorityIoService priority_service(io_context_); // 提交Type2高优先级任务(优先级设为10) priority_service.post(10, [this] { boost::asio::spawn(token_refresh_strand_, [this](boost::asio::yield_context yield) { // 令牌刷新逻辑 }); }); // 提交Type1/Type3普通优先级任务(优先级设为1) priority_service.post(1, [this] { boost::asio::spawn(io_context_, [this](boost::asio::yield_context yield) { // Type1/Type3任务逻辑 }); });
方案2:使用Boost 1.70+实验性优先级队列
从Boost 1.70开始,可使用experimental::priority_queue简化实现:
#include <boost/asio/experimental/priority_queue.hpp> // 定义优先级,数字越大优先级越高 constexpr int HIGH_PRIORITY = 10; constexpr int NORMAL_PRIORITY = 1; // 创建优先级执行器,采用大顶堆排序 auto priority_executor = boost::asio::experimental::priority_queue( io_context_.get_executor(), [](int a, int b) { return a > b; } ); // 提交Type2高优先级任务 boost::asio::post(priority_executor.wrap(HIGH_PRIORITY), [this] { boost::asio::spawn(token_refresh_strand_, [this](boost::asio::yield_context yield) { // 令牌刷新逻辑 }); }); // 提交Type1/Type3普通优先级任务 boost::asio::post(priority_executor.wrap(NORMAL_PRIORITY), [this] { boost::asio::spawn(io_context_, [this](boost::asio::yield_context yield) { // Type1/Type3任务逻辑 }); });
注:experimental模块接口可能随版本调整,需确认Boost版本兼容性。
内容的提问来源于stack exchange,提问作者Zohar81

