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)
问题原因与解决方案
核心问题
- 缓冲区数据残留:Boost.Beast的
flat_buffer是累积式缓冲区,每次async_read会将新数据追加到缓冲区末尾。你在Lambda回调中先调用on_message_handler处理数据,再清理缓冲区,但下一次读取的新数据会和未清理的旧数据拼接,形成两个JSON对象连在一起的非法格式,导致解析时出现多余的{。 - 重复注册读取回调:在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
相关产品推荐
相关产品推荐

