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

async_read的Lambda回调中调用JSON解析函数出现解析错误问题

问题:Boost.Beast开发Binance WebSocket客户端JSON解析错误

基于Boost.Beast开发Binance WebSocket客户端时,on_message函数内直接调用on_message_handler(内含nlohmann JSON解析逻辑)可正常执行;但将该handler调用放在async_read的Lambda回调函数中时,触发nlohmann::detail::parse_error,错误提示为:

[json.exception.parse_error.101] 解析第1行第687列时出错:解析值时出现语法错误 - 意外的'{';预期输入结束。

相关代码

#define BINANCE_HANDLER(f) beast::bind_front_handler(&binanceWS<A>::f, this->shared_from_this())

template <typename A> 
class binanceWS : public std::enable_shared_from_this<binanceWS<A>> {
    tcp::resolver      resolver_;
    Stream             ws_;
    beast::flat_buffer buffer_;
    std::string        host_;
    std::string        message_text_;

    std::string           wsTarget_ = "/ws/";
    char const*           host      = "stream.binance.com";
    char const*           port      = "9443";
    SPSCQueue<A>&         diff_messages_queue;
    std::function<void()> on_message_handler;

  public:
    binanceWS(net::any_io_executor ex, ssl::context& ctx, SPSCQueue<A>& q)
        : resolver_(ex)
        , ws_(ex, ctx)
        , diff_messages_queue(q) {}

    void run(char const* host, char const* port, json message, const std::string& streamName) {
        if (!SSL_set_tlsext_host_name(ws_.next_layer().native_handle(), host)) {
            throw boost::system::system_error(
                error_code(::ERR_get_error(), net::error::get_ssl_category()));
        }

        host_         = host;
        message_text_ = message.dump();
        wsTarget_ += streamName;

        resolver_.async_resolve(host_, port, BINANCE_HANDLER(on_resolve));
    }

    void on_resolve(beast::error_code ec, tcp::resolver::results_type results) {
        if (ec)
            return fail_ws(ec, "resolve");

        if (!SSL_set_tlsext_host_name(ws_.next_layer().native_handle(), host_.c_str())) {
            throw beast::system_error{
                error_code(::ERR_get_error(), net::error::get_ssl_category())};
        }

        get_lowest_layer(ws_).expires_after(30s);

        beast::get_lowest_layer(ws_).async_connect(results, BINANCE_HANDLER(on_connect));
    }

    void on_connect(beast::error_code                                           ec,
                    [[maybe_unused]] tcp::resolver::results_type::endpoint_type ep) {
        if (ec)
            return fail_ws(ec, "connect");

        ws_.next_layer().async_handshake(ssl::stream_base::client, BINANCE_HANDLER(on_ssl_handshake));
    }

    void on_ssl_handshake(beast::error_code ec) {
        if (ec)
            return fail_ws(ec, "ssl_handshake");

        beast::get_lowest_layer(ws_).expires_never();

        ws_.set_option(websocket::stream_base::timeout::suggested(beast::role_type::client));

        ws_.set_option(websocket::stream_base::decorator([](websocket::request_type& req) {
            req.set(http::field::user_agent,
                    std::string(BOOST_BEAST_VERSION_STRING) + " websocket-client-async");
        }));

        std::cout << "using host_: " << host_ << std::endl;
        ws_.async_handshake(host_, wsTarget_, BINANCE_HANDLER(on_handshake));
    }

    void on_handshake(beast::error_code ec) {
        if (ec) {
            return fail_ws(ec, "handshake");
        }

        std::cout << "Sending : " << message_text_ << std::endl;

        ws_.async_write(net::buffer(message_text_), BINANCE_HANDLER(on_write));
    }

    void on_write(beast::error_code ec, size_t bytes_transferred) {
        boost::ignore_unused(bytes_transferred);

        if (ec)
            return fail_ws(ec, "write");

        ws_.async_read(buffer_, BINANCE_HANDLER(on_message));
    }

    void on_message(beast::error_code ec, size_t bytes_transferred) {
        boost::ignore_unused(bytes_transferred);
        if (ec)
            return fail_ws(ec, "read");

       on_message_handler(); // WORKS FINE!!!

        ws_.async_read(buffer_, [this](beast::error_code ec, size_t n) {
            if (ec)
                return fail_ws(ec, "read");

            on_message_handler(); // DOESN'T WORK  
            buffer_.clear();
            ws_.async_read(buffer_, BINANCE_HANDLER(on_message));
        });
    }
    
    void subscribe_orderbook_diffs(const std::string action,const std::string symbol,short int depth_levels)
    {
        std::string stream = symbol+"@"+"depth"+std::to_string(depth_levels);

        
        on_message_handler = [this]() {
            std::cout << "Orderbook Levels Update" << std::endl;
            json payload = json::parse(beast::buffers_to_string(buffer_.cdata()));
            std::cout << payload << std::endl;
             
        };
        
        json jv = {
            { "method", action },
            { "params", {stream} },
            { "id", 1 }
        };
        run(host, port,jv, stream);
    }

};

int main() {
    net::io_context ioc;
    ssl::context    ctx{ssl::context::tlsv12_client};

    ctx.set_verify_mode(ssl::verify_peer);
    ctx.set_default_verify_paths();
    int         levels = 10;
    std::string symbol = "btcusdt";

    auto binancews = std::make_shared<binanceWS>(make_strand(ioc), ctx);
    binancews->subscribe_orderbook_diffs("SUBSCRIBE", symbol, levels);
    ioc.run();
}

错误输出

Orderbook Levels Update
terminate called after throwing an instance of 'nlohmann::detail::parse_error'
  what():  [json.exception.parse_error.101] parse error at line 1, column 687: syntax error while parsing value - unexpected '{'; expected end of input
Aborted (core dumped)

问题原因与解决方案

核心问题

  1. 缓冲区数据残留:Boost.Beast的flat_buffer是累积式缓冲区,每次async_read会将新数据追加到缓冲区末尾。你在Lambda回调中先调用on_message_handler处理数据,再清理缓冲区,但下一次读取的新数据会和未清理的旧数据拼接,形成两个JSON对象连在一起的非法格式,导致解析时出现多余的{。
  2. 重复注册读取回调:在Lambda回调中再次调用BINANCE_HANDLER(on_message),而on_message函数本身又会注册一个新的Lambda回调,导致多个异步读取操作同时往缓冲区写数据,彻底打乱数据结构。

修复代码

方案一:保持统一的回调链式调用(推荐)

修改on_message函数,处理完数据后立即清理缓冲区,然后直接注册下一次异步读取:

void on_message(beast::error_code ec, size_t bytes_transferred) {
    boost::ignore_unused(bytes_transferred);
    if (ec)
        return fail_ws(ec, "read");

    on_message_handler(); // 处理当前消息
    buffer_.clear();      // 立即清理缓冲区,避免残留

    // 注册下一次异步读取,保持单一回调链条
    ws_.async_read(buffer_, BINANCE_HANDLER(on_message));
}

方案二:修正Lambda回调逻辑(不推荐,易混乱)

如果一定要使用Lambda回调,调整操作顺序并避免重复注册:

void on_message(beast::error_code ec, size_t bytes_transferred) {
    boost::ignore_unused(bytes_transferred);
    if (ec)
        return fail_ws(ec, "read");

    ws_.async_read(buffer_, [this](beast::error_code ec, size_t n) {
        if (ec)
            return fail_ws(ec, "read");

        on_message_handler(); 
        buffer_.clear(); // 处理后立即清理缓冲区
        // 继续注册下一次读取,不要调用on_message避免嵌套
        ws_.async_read(buffer_, std::move(*this));
    });
}

额外注意事项

  • 始终在处理完当前缓冲区的数据后立即调用buffer_.clear(),确保下一次读取的是干净的新数据。
  • 避免混合使用不同的回调注册方式(BINANCE_HANDLER和Lambda),保持逻辑一致性,防止重复触发读取操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 21:55:44