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

基于Boost.Beast的WebSocket客户端与服务端实现示例求助

Boost.Beast WebSocket 客户端与服务端完整实现

一、服务端实现(含消息响应、订阅推送)

服务端实现了消息响应、连接管理、主题订阅推送三个核心功能,修复了原代码中缺少订阅逻辑、错误处理不完善的问题:

#include <boost/asio.hpp>
#include <boost/beast.hpp>
#include <nlohmann/json.hpp>
#include <iostream>
#include <memory>
#include <string>
#include <unordered_map>
#include <set>
#include <mutex>
#include <chrono>

namespace asio = boost::asio;
namespace beast = boost::beast;
namespace websocket = beast::websocket;
using tcp = asio::ip::tcp;
using json = nlohmann::json;

// 会话管理器:管理所有连接与订阅关系
class SessionManager {
public:
    void add_session(std::shared_ptr<class WebSocketSession> session) {
        std::lock_guard<std::mutex> lock(mtx_);
        sessions_.insert(session);
    }

    void remove_session(std::shared_ptr<class WebSocketSession> session) {
        std::lock_guard<std::mutex> lock(mtx_);
        sessions_.erase(session);
        // 清理该会话的所有订阅
        for (auto& [topic, subs] : subscriptions_) {
            subs.erase(session);
        }
    }

    void subscribe(std::shared_ptr<class WebSocketSession> session, const std::string& topic) {
        std::lock_guard<std::mutex> lock(mtx_);
        subscriptions_[topic].insert(session);
    }

    void unsubscribe(std::shared_ptr<class WebSocketSession> session, const std::string& topic) {
        std::lock_guard<std::mutex> lock(mtx_);
        if (subscriptions_.count(topic)) {
            subscriptions_[topic].erase(session);
        }
    }

    // 向指定主题推送消息
    void publish(const std::string& topic, const std::string& message) {
        std::lock_guard<std::mutex> lock(mtx_);
        if (!subscriptions_.count(topic)) return;

        for (auto session : subscriptions_[topic]) {
            session->send(message);
        }
    }

private:
    std::mutex mtx_;
    std::set<std::shared_ptr<class WebSocketSession>> sessions_;
    std::unordered_map<std::string, std::set<std::shared_ptr<class WebSocketSession>>> subscriptions_;
};

class WebSocketSession : public std::enable_shared_from_this<WebSocketSession> {
public:
    WebSocketSession(tcp::socket socket, SessionManager& manager)
        : ws_(std::move(socket)), manager_(manager) {
        manager_.add_session(shared_from_this());
    }

    ~WebSocketSession() {
        manager_.remove_session(shared_from_this());
    }

    void run() {
        ws_.async_accept(
            beast::bind_front_handler(
                &WebSocketSession::on_accept,
                shared_from_this()
            )
        );
    }

    // 异步发送消息,保证线程安全
    void send(const std::string& message) {
        auto self = shared_from_this();
        asio::post(ws_.get_executor(), [self, message]() {
            bool write_in_progress = !write_queue_.empty();
            write_queue_.push_back(message);
            if (!write_in_progress) {
                do_write();
            }
        });
    }

private:
    void on_accept(beast::error_code ec) {
        if (ec) {
            std::cerr << "连接接收失败: " << ec.message() << std::endl;
            return;
        }
        do_read();
    }

    void do_read() {
        ws_.async_read(
            buffer_,
            beast::bind_front_handler(
                &WebSocketSession::on_read,
                shared_from_this()
            )
        );
    }

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

        if (ec) {
            if (ec == websocket::error::closed) return;
            std::cerr << "读取消息失败: " << ec.message() << std::endl;
            return;
        }

        try {
            auto msg_str = beast::buffers_to_string(buffer_.data());
            json msg = json::parse(msg_str);
            buffer_.consume(buffer_.size());

            // 处理不同类型请求
            if (msg.contains("echo")) {
                json response = {
                    {"type", "echo_response"},
                    {"original", msg["echo"]}
                };
                send(response.dump());
            } else if (msg.contains("subscribe")) {
                std::string topic = msg["subscribe"];
                manager_.subscribe(shared_from_this(), topic);
                json response = {
                    {"type", "subscribe_ack"},
                    {"topic", topic},
                    {"status", "success"}
                };
                send(response.dump());
            } else if (msg.contains("unsubscribe")) {
                std::string topic = msg["unsubscribe"];
                manager_.unsubscribe(shared_from_this(), topic);
                json response = {
                    {"type", "unsubscribe_ack"},
                    {"topic", topic},
                    {"status", "success"}
                };
                send(response.dump());
            } else {
                json response = {
                    {"type", "error"},
                    {"message", "未知消息类型"}
                };
                send(response.dump());
            }

            do_read();
        } catch (const std::exception& e) {
            std::cerr << "消息处理失败: " << e.what() << std::endl;
            json error_resp = {{"type", "error"}, {"message", "消息格式错误"}};
            send(error_resp.dump());
            do_read();
        }
    }

    void do_write() {
        ws_.text(ws_.got_text());
        ws_.async_write(
            asio::buffer(write_queue_.front()),
            beast::bind_front_handler(
                &WebSocketSession::on_write,
                shared_from_this()
            )
        );
    }

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

        if (ec) {
            std::cerr << "发送消息失败: " << ec.message() << std::endl;
            return;
        }

        write_queue_.pop_front();
        if (!write_queue_.empty()) {
            do_write();
        }
    }

    websocket::stream<tcp::socket> ws_;
    beast::flat_buffer buffer_;
    SessionManager& manager_;
    std::deque<std::string> write_queue_;
};

class WebSocketServer {
public:
    WebSocketServer(asio::io_context& ioc, tcp::endpoint endpoint)
        : acceptor_(ioc, endpoint), manager_(std::make_shared<SessionManager>()) {
        do_accept();

        // 模拟定时向news主题推送消息
        std::thread push_thread([this]() {
            int count = 0;
            while (true) {
                std::this_thread::sleep_for(std::chrono::seconds(5));
                json push_msg = {
                    {"type", "push"},
                    {"topic", "news"},
                    {"content", "定时推送消息 " + std::to_string(++count)}
                };
                manager_->publish("news", push_msg.dump());
                std::cout << "已向news主题推送消息" << std::endl;
            }
        });
        push_thread.detach();
    }

private:
    void do_accept() {
        acceptor_.async_accept(
            beast::bind_front_handler(
                &WebSocketServer::on_accept,
                this
            )
        );
    }

    void on_accept(beast::error_code ec, tcp::socket socket) {
        if (ec) {
            std::cerr << "接受连接失败: " << ec.message() << std::endl;
        } else {
            std::make_shared<WebSocketSession>(std::move(socket), *manager_)->run();
        }
        do_accept();
    }

    tcp::acceptor acceptor_;
    std::shared_ptr<SessionManager> manager_;
};

int main() {
    std::cout << "WebSocket服务端已启动,监听端口9002..." << std::endl;
    try {
        asio::io_context ioc{1};
        tcp::endpoint endpoint(tcp::v4(), 9002);
        WebSocketServer server(ioc, endpoint);
        ioc.run();
    } catch (const std::exception& e) {
        std::cerr << "服务端启动失败: " << e.what() << std::endl;
        return EXIT_FAILURE;
    }
    return EXIT_SUCCESS;
}

服务端核心说明

  • 会话管理:通过SessionManager维护所有连接,线程安全地处理订阅、取消订阅操作。
  • 消息响应:解析JSON消息,对合法请求返回对应响应,非法消息返回错误提示,避免客户端无响应等待。
  • 订阅推送:模拟定时向指定主题推送消息,遍历订阅该主题的所有会话异步发送。
  • 异步发送队列:每个会话维护消息队列,保证异步发送的顺序,避免并发发送冲突。

二、客户端实现(含订阅、消息收发)

修复了原代码中多线程直接操作WebSocket的线程安全问题,使用asio::strand序列化操作,实现完整的订阅、收发逻辑:

#include <boost/asio.hpp>
#include <boost/beast.hpp>
#include <nlohmann/json.hpp>
#include <iostream>
#include <memory>
#include <string>
#include <thread>

namespace asio = boost::asio;
namespace beast = boost::beast;
namespace websocket = beast::websocket;
using tcp = asio::ip::tcp;
using json = nlohmann::json;

class WebSocketClient : public std::enable_shared_from_this<WebSocketClient> {
public:
    WebSocketClient(asio::io_context& ioc)
        : ioc_(ioc), strand_(asio::make_strand(ioc)) {}

    void connect(const std::string& host, const std::string& port) {
        resolver_.async_resolve(host, port,
            asio::bind_executor(strand_,
                beast::bind_front_handler(
                    &WebSocketClient::on_resolve,
                    shared_from_this()
                )
            )
        );
    }

    // 线程安全的消息发送
    void send(const json& msg) {
        asio::post(strand_, [this, msg_str = msg.dump()]() {
            bool write_in_progress = !write_queue_.empty();
            write_queue_.push_back(msg_str);
            if (!write_in_progress) {
                do_write();
            }
        });
    }

private:
    void on_resolve(beast::error_code ec, tcp::resolver::results_type results) {
        if (ec) {
            std::cerr << "解析地址失败: " << ec.message() << std::endl;
            return;
        }

        asio::async_connect(ws_.next_layer(), results.begin(), results.end(),
            asio::bind_executor(strand_,
                beast::bind_front_handler(
                    &WebSocketClient::on_connect,
                    shared_from_this()
                )
            )
        );
    }

    void on_connect(beast::error_code ec, tcp::resolver::results_type::endpoint_type) {
        if (ec) {
            std::cerr << "连接服务端失败: " << ec.message() << std::endl;
            return;
        }

        ws_.async_handshake("127.0.0.1", "/",
            asio::bind_executor(strand_,
                beast::bind_front_handler(
                    &WebSocketClient::on_handshake,
                    shared_from_this()
                )
            )
        );
    }

    void on_handshake(beast::error_code ec) {
        if (ec) {
            std::cerr << "握手失败: " << ec.message() << std::endl;
            return;
        }

        std::cout << "已连接到服务端" << std::endl;
        do_read();
    }

    void do_read() {
        ws_.async_read(buffer_,
            asio::bind_executor(strand_,
                beast::bind_front_handler(
                    &WebSocketClient::on_read,
                    shared_from_this()
                )
            )
        );
    }

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

        if (ec) {
            if (ec == websocket::error::closed) {
                std::cerr << "连接已关闭" << std::endl;
            } else {
                std::cerr << "读取消息失败: " << ec.message() << std::endl;
            }
            return;
        }

        try {
            auto msg_str = beast::buffers_to_string(buffer_.data());
            json msg = json::parse(msg_str);
            buffer_.consume(buffer_.size());

            // 处理服务端消息
            if (msg["type"] == "echo_response") {
                std::cout << "\n收到echo响应: " << msg["original"] << std::endl;
            } else if (msg["type"] == "subscribe_ack") {
                std::cout << "\n订阅主题" << msg["topic"] << "成功" << std::endl;
            } else if (msg["type"] == "unsubscribe_ack") {
                std::cout << "\n取消订阅主题" << msg["topic"] << "成功" << std::endl;
            } else if (msg["type"] == "push") {
                std::cout << "\n收到推送消息(" << msg["topic"] << "): " << msg["content"] << std::endl;
            } else if (msg["type"] == "error") {
                std::cout << "\n收到错误消息: " << msg["message"] << std::endl;
            }

            do_read();
        } catch (const std::exception& e) {
            std::cerr << "解析消息失败: " << e.what() << std::endl;
            do_read();
        }
    }

    void do_write() {
        ws_.text(ws_.got_text());
        ws_.async_write(asio::buffer(write_queue_.front()),
            asio::bind_executor(strand_,
                beast::bind_front_handler(
                    &WebSocketClient::on_write,
                    shared_from_this()
                )
            )
        );
    }

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

        if (ec) {
            std::cerr << "发送消息失败: " << ec.message() << std::endl;
            return;
        }

        write_queue_.pop_front();
        if (!write_queue_.empty()) {
            do_write();
        }
    }

    asio::io_context& ioc_;
    asio::strand<asio::io_context::executor_type> strand_;
    tcp::resolver resolver_{strand_};
    websocket::stream<tcp::socket> ws_{strand_};
    beast::flat_buffer buffer_;
    std::deque<std::string> write_queue_;
};

int main() {
    try {
        asio::io_context ioc;
        auto client = std::make_shared<WebSocketClient>(ioc);
        client->connect("127.0.0.1", "9002");

        // 启动io线程
        std::thread io_thread([&ioc]() {
            ioc.run();
        });

        // 控制台输入处理
        std::string input;
        while (true) {
            std::cout << "\n输入指令(echo <内容>/subscribe <主题>/unsubscribe <主题>/stop): ";
            std::getline(std::cin, input);

            if (input == "stop") break;
            else if (input.substr(0,4) == "echo" && input.size()>5) {
                client->send({{"echo", input.substr(5)}});
            } else if (input.substr(0,9) == "subscribe" && input.size()>10) {
                client->send({{"subscribe", input.substr(10)}});
            } else if (input.substr(0,11) == "unsubscribe" && input.size()>12) {
                client->send({{"unsubscribe", input.substr(12)}});
            } else {
                std::cout << "指令格式错误,请重新输入" << std::endl;
            }
        }

        // 优雅关闭连接
        asio::post(client->ws_.get_executor(), [client]() {
            client->ws_.close(websocket::close_code::normal);
        });
        io_thread.join();
    } catch (const std::exception& e) {
        std::cerr << "客户端运行错误: " << e.what() << std::endl;
        return EXIT_FAILURE;
    }
    return EXIT_SUCCESS;
}

客户端核心说明

  • 线程安全:通过
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:48:51