You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Boost.Asio与Boost.Interprocess Message Queue的服务器异步读写阻塞问题求解

解决Boost Asio异步读阻塞导致无法处理消息队列的问题

你的问题核心在于Asio的事件循环被TCP异步读操作占据,同时没有持续监听消息队列的新数据,导致服务器没法及时把消息队列里的响应写回TCP连接。我来帮你拆解问题并给出可行的解决方案。

问题分析

先看你的代码里几个关键的问题点:

  • do_write只执行一次:在start()里调用了一次do_write(),但如果消息队列一开始没有数据,try_receive直接返回,之后就再也不会主动检查消息队列了。即使后面消息队列有新的响应数据,也没有触发写操作的逻辑。
  • 异步读的回调逻辑有缺陷:do_read_header()的回调里不管有没有错误(比如socket断开),都无条件调用do_read_header(),这会导致无效循环,甚至占用strand的执行时间,影响其他操作。
  • 消息队列与Asio事件循环脱节:Boost进程间消息队列的操作不属于Asio的IO事件,Asio的io_context不会主动感知到消息队列的新数据,必须我们主动轮询或用线程监听。

可行解决方案

这里提供两种常用思路,你可以根据业务场景选择:

方案1:用Asio定时器定时轮询消息队列

这种方案不需要额外线程,通过定时器定期检查消息队列,适合消息队列数据不是特别频繁的场景。

修改你的Session类,关键改动如下:

class Session : public std::enable_shared_from_this<Session> {
public:
    Session(tcp::socket socket)
        : socket_(std::move(socket))
        , response_queue_(boost::interprocess::open_or_create, "response_queue", 100, 100)
        , request_queue_(boost::interprocess::open_or_create, "request_queue", 100, 100)
        , timer_(socket_.get_executor()) // 复用socket的执行器初始化定时器
    { }

    void start() {
        post(strand_, [this, self = shared_from_this()] {
            start_response_poll(); // 启动消息队列轮询
            do_read_header();
        });
    }

private:
    void start_response_poll() {
        auto self(shared_from_this());
        // 每隔100ms轮询一次(时间间隔可根据需求调整)
        timer_.expires_after(std::chrono::milliseconds(100));
        timer_.async_wait(strand_.wrap([this, self](boost::system::error_code ec) {
            if (!ec && socket_.is_open()) {
                do_write();
                start_response_poll(); // 继续下一轮轮询
            }
        }));
    }

    void do_write() {
        auto self(shared_from_this());
        Message msg;
        boost::interprocess::message_queue::size_type recvd_size;
        unsigned int priority = 0;
        
        // 有数据才处理写操作
        if(response_queue_.try_receive(msg.body(), msg.max_body_length, recvd_size, priority)) {
            try {
                msg.body_length(recvd_size);
                msg.encode_header();
                std::memcpy(data_, msg.data(), msg.length());
            } catch (std::exception& e ) {
                std::cout << e.what();
                return;
            }
            std::cout << "write " << msg.body() << "\n";
            boost::asio::async_write(socket_,
                boost::asio::buffer(data_, msg.length()),
                strand_.wrap([this, self](boost::system::error_code ec, std::size_t /*length*/) {
                    if (ec) {
                        std::cerr << "write error:" << ec.value() << " message: " << ec.message() << "\n";
                        socket_.close();
                        timer_.cancel(); // 关闭定时器
                    }
                }));
        }
    }

    // 修复异步读头部的回调逻辑
    void do_read_header() {
        auto self(shared_from_this());
        std::cout << "do_read_header\n";
        boost::asio::async_read(socket_,
            boost::asio::buffer(res.data(), res.header_length),
            strand_.wrap([this, self](boost::system::error_code ec, std::size_t /*length*/) {
                if (ec) {
                    std::cerr << "read header error:" << ec.value() << " message: " << ec.message() << "\n";
                    socket_.close();
                    timer_.cancel();
                    return;
                }
                if (res.decode_header()) {
                    do_read_body();
                } else {
                    socket_.close();
                    timer_.cancel();
                }
            }));
    }

    // 修复异步读body的回调逻辑
    void do_read_body() {
        auto self(shared_from_this());
        std::cout << "do_read_body\n";
        boost::asio::async_read(socket_,
            boost::asio::buffer(res.body(), res.body_length()),
            strand_.wrap([this, self](boost::system::error_code ec, std::size_t length) {
                if (ec) {
                    std::cerr << "read body error:" << ec.value() << " message: " << ec.message() << "\n";
                    socket_.close();
                    timer_.cancel();
                    return;
                }
                if (!length) {
                    socket_.close();
                    timer_.cancel();
                    return;
                }
                if constexpr (log_active) {
                    request_log << "read " << res.body() << "\n";
                    request_log.flush();
                }
                std::cout << "read " << res.body() << "\n";
                request_queue_.send(res.body(), res.body_length(), 0);
                do_read_header(); // 继续读取下一个消息
            }));
    }

    // ... 原有成员
    boost::asio::steady_timer timer_; // 添加定时器成员
};

方案2:单独线程监听消息队列

如果消息队列的数据非常频繁,定时器轮询可能有延迟,这时候可以创建单独线程阻塞监听消息队列,一旦有数据就post到strand执行写操作。

关键代码修改如下:

class Session : public std::enable_shared_from_this<Session> {
public:
    Session(tcp::socket socket)
        : socket_(std::move(socket))
        , response_queue_(boost::interprocess::open_or_create, "response_queue", 100, 100)
        , request_queue_(boost::interprocess::open_or_create, "request_queue", 100, 100)
        , response_listener_(&Session::listen_response_queue, this) // 启动监听线程
    { }

    ~Session() {
        // 清理资源
        socket_.close();
        response_queue_.try_send("", 0, 0); // 唤醒阻塞的receive
        if (response_listener_.joinable()) {
            response_listener_.join();
        }
    }

    void start() {
        post(strand_, [this, self = shared_from_this()] {
            do_read_header();
        });
    }

private:
    // 监听消息队列的线程函数
    void listen_response_queue() {
        while (socket_.is_open()) {
            Message msg;
            boost::interprocess::message_queue::size_type recvd_size;
            unsigned int priority = 0;
            try {
                response_queue_.receive(msg.body(), msg.max_body_length, recvd_size, priority); // 阻塞等待数据
            } catch (const boost::interprocess::interprocess_exception& e) {
                std::cerr << "message queue receive error: " << e.what() << "\n";
                break;
            }
            if (!socket_.is_open()) break;
            // 把写操作post到strand,保证线程安全
            post(strand_, [this, self = shared_from_this(), msg = std::move(msg), recvd_size]() mutable {
                do_write(msg, recvd_size);
            });
        }
    }

    void do_write(Message msg, boost::interprocess::message_queue::size_type recvd_size) {
        try {
            msg.body_length(recvd_size);
            msg.encode_header();
            std::memcpy(data_, msg.data(), msg.length());
        } catch (std::exception& e ) {
            std::cout << e.what();
            return;
        }
        std::cout << "write " << msg.body() << "\n";
        boost::asio::async_write(socket_,
            boost::asio::buffer(data_, msg.length()),
            strand_.wrap([this](boost::system::error_code ec, std::size_t /*length*/) {
                if (ec) {
                    std::cerr << "write error:" << ec.value() << " message: " << ec.message() << "\n";
                    socket_.close();
                }
            }));
    }

    // ... 原有成员(do_read_header和do_read_body的修复同方案1)
    std::thread response_listener_; // 添加监听线程成员
};

额外注意事项

  • strand的正确使用:所有操作socket的代码都要通过strand执行,避免线程安全问题,上面的代码已经用strand_.wrap或post(strand_, ...)保证。
  • 资源清理:当socket关闭时,要及时停止定时器或join监听线程,避免资源泄漏。
  • 异常处理:消息队列的receive/try_receive可能抛出interprocess_exception,要做好捕获处理。

内容的提问来源于stack exchange,提问作者yaodav

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 11:27:56