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

WebSocket等待接收报价时发消息崩溃,如何终止on_read回调?

问题描述

基于Boost Beast异步SSL客户端示例编写代码,目标是注册WebSocket服务获取实时公司报价,每次收到报价后调用async_read继续接收后续报价。目前遇到两个问题:

  1. 等待小公司报价可能耗时数小时,程序会处于等待状态,无法发送其他公司的报价请求;
  2. 尝试用post函数在正确上下文线程调用async_write发送新订阅请求时,程序直接崩溃。

请问有没有办法强制完成on_read回调,以获得发送新消息的机会?

简化后的代码(未含互斥锁)
void
on_read(
    beast::error_code ec,
    std::size_t bytes_transferred)
{
    boost::ignore_unused(bytes_transferred);

    if(ec)
        return fail2(ec, "read");

    
    std::string mycontent = beast::buffers_to_string(buffer_.data());
    cout << mycontent << endl;
    buffer_.clear();
    
    ws_.async_read(
        buffer_,
        beast::bind_front_handler(
            &session::on_read,
            shared_from_this()));
}

void subscribe(const std::string &symbol)
{
    // 保存消息到队列
    std::string text = "{\"action\": \"subscribe\", \"symbols\": \"" + symbol + "\"}";
    msgqueue_.push_back(text);
    boost::asio::post(ioc_, beast::bind_front_handler(&session::_subscription_to_post, shared_from_this()));
}

void _subscription_to_post()
{
    if (msgqueue_.empty())
        return;

    // 发送消息
    ws_.async_write(
        net::buffer(msgqueue_.front()),
        beast::bind_front_handler(
            &session::on_write,
            shared_from_this()));
    msgqueue_.pop_front();
}
问题分析与解决

1. 异步写操作崩溃的原因及修复

你的代码存在两个核心问题导致崩溃:

  • 提前销毁缓冲区引用的内存:调用async_write后立刻执行msgqueue_.pop_front(),但异步写操作不会立即完成,此时缓冲区引用的字符串已被销毁,后续异步操作访问失效内存会触发崩溃;
  • 未处理并发写操作:Boost Beast的WebSocket不允许同时存在多个未完成的异步写操作,直接发起新写操作会导致未定义行为。

修复方案:

// 在session类中添加成员变量
bool is_writing_ = false;
std::mutex queue_mutex_;

void on_write(beast::error_code ec, std::size_t bytes_transferred)
{
    boost::ignore_unused(bytes_transferred);
    if(ec)
        return fail2(ec, "write");

    is_writing_ = false;
    // 写操作完成后再移除队列元素
    {
        std::lock_guard<std::mutex> lock(queue_mutex_);
        if(!msgqueue_.empty())
            msgqueue_.pop_front();
    }

    // 队列还有消息则继续发送
    if(!msgqueue_.empty())
        _subscription_to_post();
}

void subscribe(const std::string &symbol)
{
    std::string text = "{\"action\": \"subscribe\", \"symbols\": \"" + symbol + "\"}";
    {
        std::lock_guard<std::mutex> lock(queue_mutex_);
        msgqueue_.push_back(text);
    }
    boost::asio::post(ioc_, beast::bind_front_handler(&session::_subscription_to_post, shared_from_this()));
}

void _subscription_to_post()
{
    std::string msg;
    {
        std::lock_guard<std::mutex> lock(queue_mutex_);
        if(msgqueue_.empty())
            return;
        msg = msgqueue_.front();
    }

    if(is_writing_)
        return;

    is_writing_ = true;
    ws_.async_write(
        net::buffer(msg),
        beast::bind_front_handler(&session::on_write, shared_from_this()));
}

同时,给队列操作添加互斥锁,避免多线程调用subscribe时的并发访问问题。

2. 等待报价时无法发送新请求的解决

你不需要强制完成on_read回调——Boost Beast的异步模型本身是非阻塞的,async_read只是发起一个异步等待操作,不会阻塞IO线程,线程会继续处理post过来的_subscription_to_post任务。

之前误以为程序阻塞,大概率是崩溃问题导致的异常状态,修复崩溃后,即使async_read在等待报价,你依然可以通过subscribe提交新请求,IO线程会及时处理并发送新的订阅消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 11:05:15