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

Boost ASIO async_read_until无法连续读取消息问题排查

问题描述

我正在用Boost.Asio实现一个网络库,下面是TCP类的读取函数。测试时发现,客户端和服务器完成握手后,服务器连续发两个带分隔符的消息,但客户端只能读到第一个;等发第三个消息时,才会读到之前漏掉的第二个。bytes变量只表示单个消息的长度。代码里有没有明显问题?(已移除无关代码)

读取函数代码

boost::asio::streambuf _recBuf;
ssl::stream<ip::tcp::socket&> _stream;

void SomeClass::read(const String &delim) {
    boost::asio::async_read_until(_stream, _recBuf, delim,
    boost::asio::bind_executor(_readStrand,
    [this, delim, self = shared_from_this()] (const boost::system::error_code& ec, const std::size_t bytes) {
        if (!ec) {

            if (bytes <= delim.size()) {
                return;
            }

            String data(boost::asio::buffers_begin(_recBuf.data()), boost::asio::buffers_begin(_recBuf.data()) + (bytes - delim.size()));
            _recBuf.consume(bytes);

            read(delim);
        }
        else if (ec == system::errc::operation_canceled) {
            
        }
        else {
            
        }
    }));
}

我试过用commit()方法和istream修改读取_recBuf的方式,调用_stream.lowest_layer().available()查看连续两个消息的可用字节数,读取第一个消息后返回0。

补充:写入方法代码

会不会是写入方法有问题?

void TCPSBase::write(const Message &msg) {
String data = msg.toJsonStr() + '^';
asio::async_write(_stream, asio::buffer(data),
asio::bind_executor(_writeStrand,
[this, data, self = shared_from_this()] (const boost::system::error_code& ec, const std::size_t bytes) {
    LogHandler lh(Logger::getInstance(), "TCPSBase::write::lambda");
    if (!ec) {
        lh.d("wrote " + std::to_string(bytes) + " bytes");
        lh.d(data);
    }
    else if (ec == system::errc::operation_canceled) {
        lh.d("canceled");
    }
    else {
        lh.e("disconnecting after error during write: " + ec.message());
        stop();
    }
}));

问题分析与解决方案

你的问题核心出在读取回调的异常分支处理和streambuf的读取逻辑上,具体问题点和修复方案如下:

1. 致命问题:异常分支中断读取流程

代码里if (bytes <= delim.size())分支直接return,没有再次调用read(delim)发起下一次读取。如果某种场景触发这个分支(比如网络分包导致缓冲区里只剩分隔符),读取流程会直接中断,后续缓冲区里的剩余消息(比如你测试中的第二个消息)不会被处理,直到新数据到来(第三个消息)才会触发async_read_until,此时才会一并处理之前遗留的消息。

修复:
无论是否触发这个分支,处理后都要继续发起读取:

if (bytes <= delim.size()) {
    _recBuf.consume(bytes); // 先丢弃无效的分隔符数据
    read(delim); // 必须继续发起下一次读取
    return;
}

2. 潜在问题:直接操作缓冲区迭代器的风险

你直接用buffers_begin(_recBuf.data())构造字符串,streambuf内部可能由多个内存块组成,直接通过迭代器截取可能会导致数据读取不完整(比如消息跨缓冲区时)。更可靠的方式是用std::istream读取streambuf。

修复:
替换数据读取逻辑:

std::istream is(&_recBuf);
std::string data;
// 针对单字符分隔符'^'的读取方式,多字符分隔符需要额外处理
std::getline(is, data, '^');
// async_read_until返回的bytes包含分隔符长度,所以直接consume整个bytes即可
_recBuf.consume(bytes);

3. 写入方法的验证

服务器的写入方法用_writeStrand保护了异步操作,不会出现并发写入的问题,连续调用write会被strand排队执行,保证数据发送的完整性,所以写入逻辑没有明显问题。

调试建议

添加日志跟踪streambuf状态和读取流程,定位具体阻塞点:

void SomeClass::read(const String &delim) {
    // 日志:发起读取时的缓冲区大小
    LOG_DEBUG("发起async_read_until,当前缓冲区可用字节数:%zu", _recBuf.size());
    boost::asio::async_read_until(_stream, _recBuf, delim,
    boost::asio::bind_executor(_readStrand,
    [this, delim, self = shared_from_this()] (const boost::system::error_code& ec, const std::size_t bytes) {
        // 日志:回调触发时的状态
        LOG_DEBUG("读取回调触发,错误码:%d,读取字节数:%zu,缓冲区剩余字节数:%zu", 
                  ec.value(), bytes, _recBuf.size());
        if (!ec) {
            if (bytes <= delim.size()) {
                LOG_WARN("读取到无效数据,仅包含分隔符,字节数:%zu", bytes);
                _recBuf.consume(bytes);
                read(delim);
                return;
            }

            std::istream is(&_recBuf);
            std::string data;
            std::getline(is, data, '^');
            LOG_DEBUG("读取到有效数据:%s", data.c_str());
            _recBuf.consume(bytes);
            LOG_DEBUG("处理后缓冲区剩余字节数:%zu", _recBuf.size());

            read(delim);
        }
        else if (ec == system::errc::operation_canceled) {
            LOG_DEBUG("读取操作被取消");
        }
        else {
            LOG_ERROR("读取错误:%s", ec.message().c_str());
        }
    }));
}

通过日志可以清晰看到:

  • 第一次读取后缓冲区剩余的字节数是否包含第二个消息
  • 第二次async_read_until是否立即触发回调
  • 是否存在异常分支中断读取的情况

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 10:11:00