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

如何在C++类中并行实现UDP数据包收发及ACK确认机制

C++ 多线程带ACK UDP通信类实现方案修正

现有代码的核心问题

  • 你当前sendm方法中创建线程后立即调用th.join()是致命错误:join()会阻塞当前调用线程,直到发送线程退出才会返回,你的for循环调用send时会永远卡在第一个发送任务,完全达不到并行发送5个消息的效果。
  • 没有线程停止标识:C++ 不支持直接强制终止线程,直接存储std::thread对象无法安全停止线程,必须通过原子标志位通知线程主动退出。
  • 共享数据无锁保护:host_map、ACK表、线程列表都是多线程并发访问的资源,不加互斥锁会出现数据竞争,导致未定义行为。
  • 缺少序列号匹配逻辑:仅靠主机ID无法匹配同一主机发送的多条消息的ACK,必须给每条消息分配唯一序列号,ACK报文中要带回该序列号才能对应到正确的发送线程。

可行实现方案

1. 类成员变量设计

#include <atomic>
#include <mutex>
#include <unordered_map>
#include <vector>
#include <thread>
#include <memory>
#include <algorithm>

// 发送任务结构体,每个消息对应一个实例
struct SendTask {
    std::atomic<bool> stop = false; // 线程停止标志
    std::thread th;
    std::string msg;
    in_addr_t dst_ip;
    unsigned short dst_port;
    int seq; // 消息唯一序列号
};

class A {
private:
    int obj_socket;
    std::atomic<bool> is_running = true; // 全局运行标志
    std::mutex table_mtx; // 共享资源互斥锁
    std::unordered_map<int /* 序列号seq */, SendTask*> seq_to_task; // ACK匹配表
    std::vector<std::unique_ptr<SendTask>> all_tasks; // 所有发送任务列表,用于shutdown
    std::thread listen_th; // 监听线程

    int resolveId(unsigned short port, in_addr_t ip);
    int resolveSeq(); // 生成唯一序列号,可用原子自增实现
    void sendMessage(SendTask* task);
    void listen_loop();
    void receive();
public:
    A(in_addr_t my_ip, unsigned short port);
    ~A() { shutdown(); }
    void listen();
    void send(std::string msg, in_addr_t dst_ip, unsigned short dst_port);
    void shutdown();
};

2. 核心方法实现

listen方法

仅负责启动监听线程,不阻塞主线程

void A::listen() {
    listen_th = std::thread(&A::listen_loop, this);
}

void A::listen_loop() {
    while(is_running) {
        receive();
    }
}

send方法

创建发送任务,启动线程后直接返回,不阻塞

void A::send(std::string m, in_addr_t dst_ip, unsigned short dst_port) {
    std::lock_guard<std::mutex> lock(table_mtx);
    // 生成唯一序列号
    int seq = resolveSeq();
    // 创建发送任务
    auto task = std::make_unique<SendTask>();
    task->msg = std::move(m);
    task->dst_ip = dst_ip;
    task->dst_port = dst_port;
    task->seq = seq;
    // 启动发送线程
    task->th = std::thread(&A::sendMessage, this, task.get());
    // 存入映射表和全局任务列表
    seq_to_task[seq] = task.get();
    all_tasks.push_back(std::move(task));
}

sendMessage方法

通过原子标志判断是否退出循环

void A::sendMessage(SendTask* task) {
    struct sockaddr_in dst{};
    dst.sin_family = AF_INET;
    dst.sin_port = htons(task->dst_port);
    dst.sin_addr.s_addr = task->dst_ip;

    int t = 50;
    // 收到停止信号或者全局关闭时退出
    while(!task->stop && is_running) {
        // 发送的消息要带上seq序列号,方便对端返回ACK时匹配
        std::string send_buf = std::to_string(task->seq) + "|" + task->msg;
        if (sendto(obj_socket, send_buf.c_str(), send_buf.size(), 0,  
                   reinterpret_cast<const sockaddr*>(&dst), sizeof(dst)) < 0) {
            std::cerr << "Error sendto\n";
            break;
        }
        std::this_thread::sleep_for(std::chrono::milliseconds(t));
        t = std::min(t + 10, 1000); // 加个最大延迟上限,避免无限增长
    }
}

receive方法

收到ACK时停止对应发送线程并回收资源

void A::receive() {
    char buffer[1500];
    sockaddr_in from;
    socklen_t fromlen = sizeof(from);
    ssize_t tmp = recvfrom(obj_socket, buffer, 1500, 0, 
                           reinterpret_cast<sockaddr*>(&from), &fromlen);
    if (tmp < 0) {
        if (!is_running) return; // 全局关闭时直接退出
        std::cerr << "Error receive from\n";
        return;
    }
    buffer[tmp] = '\0';
    // 解析ACK报文,取出对应的seq序列号(这里要和对端约定ACK报文格式)
    int ack_seq = parse_ack_seq(buffer, tmp);
    if (ack_seq < 0) {
        // 不是ACK报文,处理普通消息然后回复ACK
        int id = resolveId(from.sin_port, from.sin_addr.s_addr);
        // 你的业务逻辑
        send_ack(from, id, get_seq_from_msg(buffer));
        return;
    }
    // 匹配到ACK,停止对应发送线程
    std::lock_guard<std::mutex> lock(table_mtx);
    auto it = seq_to_task.find(ack_seq);
    if (it != seq_to_task.end()) {
        SendTask* task = it->second;
        task->stop = true;
        if (task->th.joinable()) {
            task->th.join(); // 这里join回收线程资源
        }
        // 从映射表和全局任务列表删除
        seq_to_task.erase(it);
        all_tasks.erase(std::remove_if(all_tasks.begin(), all_tasks.end(),
            [task](const std::unique_ptr<SendTask>& t) { return t.get() == task; }),
            all_tasks.end());
    }
}

shutdown方法

停止所有线程并回收资源

void A::shutdown() {
    if (!is_running.exchange(false)) return; // 避免重复调用
    // 关闭socket,让recvfrom立即返回错误
    close(obj_socket);
    // 回收监听线程
    if (listen_th.joinable()) {
        listen_th.join();
    }
    // 停止所有发送线程并回收
    std::lock_guard<std::mutex> lock(table_mtx);
    for (auto& task : all_tasks) {
        task->stop = true;
        if (task->th.joinable()) {
            task->th.join();
        }
    }
    all_tasks.clear();
    seq_to_task.clear();
    // 写入关闭日志
    write_log("All threads stopped, shutdown complete");
}

两种方案对比

  • 你当前的单消息单线程方案:实现简单,适合并发消息量少的场景,消息量超过几十上百个时,线程切换开销会很高。
  • 发送队列方案:更高效,适合高并发场景,实现逻辑是:类内部维护一个待发送消息队列,开固定3-5个工作线程从队列取消息发送,维护每个消息的重试状态,不需要为每个消息开线程,资源消耗更低,你可以先跑通单线程版本后再优化为队列方案。

额外注意点

  • UDP的sendto是线程安全的,多个发送线程同时调用不会有问题。
  • 序列号可以用std::atomic<int>全局自增实现,保证唯一。
  • 可以给发送逻辑加最大重试次数,避免消息永远发不出去占着资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 19:36:04