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

多线程Socket发送场景下的线程安全队列容器实现需求

线程安全的Socket发送队列实现方案

核心思路:生产者-消费者模型

用线程安全队列存待发送数据,仅用单个线程负责从队列取数据并调用send(),彻底避免多线程直接操作Socket和容器的线程安全问题。

1. 基于std::queue的线程安全队列实现

没必要用std::vector,std::queue更贴合先进先出的发送需求,配合互斥锁和条件变量实现线程安全,还能避免空轮询浪费CPU:

#include <queue>
#include <mutex>
#include <condition_variable>
#include <vector>

class SafeSendQueue {
private:
    std::queue<std::vector<char>> queue_;
    std::mutex mutex_;
    std::condition_variable cv_;
    bool stop_flag_ = false;

public:
    // 生产者线程调用:添加待发送数据
    void push(std::vector<char> data) {
        std::lock_guard<std::mutex> lock(mutex_);
        queue_.push(std::move(data));
        cv_.notify_one(); // 通知发送线程有数据待处理
    }

    // 发送线程调用:阻塞获取待发送数据,返回false表示停止
    bool pop(std::vector<char>& out_data) {
        std::unique_lock<std::mutex> lock(mutex_);
        // 等待直到队列有数据或收到停止信号
        cv_.wait(lock, [this]() { return !queue_.empty() || stop_flag_; });
        
        if (stop_flag_ && queue_.empty()) {
            return false;
        }

        out_data = std::move(queue_.front());
        queue_.pop();
        return true;
    }

    // 停止发送线程
    void stop() {
        std::lock_guard<std::mutex> lock(mutex_);
        stop_flag_ = true;
        cv_.notify_all();
    }
};

2. 单发送线程的业务逻辑

启动专门的发送线程,循环从安全队列取数据并调用send():

#include <thread>
#include <sys/socket.h> // Linux环境;Windows替换为winsock2.h

void send_worker(SafeSendQueue& queue, int sock_fd) {
    std::vector<char> data;
    while (queue.pop(data)) {
        // 处理send的部分发送情况:循环发送直到全部数据发出
        ssize_t sent_total = 0;
        ssize_t remaining = data.size();
        while (remaining > 0) {
            ssize_t sent = send(sock_fd, data.data() + sent_total, remaining, 0);
            if (sent == -1) {
                // 此处处理发送错误:比如断开重连、记录日志等
                break;
            }
            sent_total += sent;
            remaining -= sent;
        }
    }
}

// 主逻辑示例
int main() {
    int sock_fd = ...; // 假设已完成Socket创建与连接
    SafeSendQueue send_queue;
    std::thread sender(send_worker, std::ref(send_queue), sock_fd);

    // 其他业务线程直接调用send_queue.push(data)添加待发送数据

    // 程序退出时清理
    send_queue.stop();
    sender.join();
    close(sock_fd);
    return 0;
}

3. 替代方案

  • 无锁队列:如果是高并发场景,可使用boost::lockfree::queue(依赖Boost库),避免互斥锁的性能开销。
  • Asio Strand:如果项目使用Boost.Asio或C++20的std::asio,直接用io_context::strand将Socket的send操作绑定到同一执行序列,Asio会自动处理任务排队,无需手动实现队列。

关键注意事项

  • 禁止多线程直接调用同一Socket的send(),即使队列安全,并发调用也会导致数据字节交织,破坏协议格式。
  • 必须处理send()的返回值:send()可能仅发送部分数据,需要循环补全发送;若返回-1,需根据错误码处理连接断开等异常。
  • 停止线程时,可根据业务需求决定是否发送完队列剩余数据,或直接丢弃。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 09:37:27