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

异步TCP服务器中N线程访问数据的线程安全问题

这个场景其实是典型的单生产者-多消费者同步问题,还带了几个关键约束:写操作排他、读操作可并发,生产者必须等所有消费者完成一轮处理才能生成下一批数据,消费者处理完要等新数据同时不能丢请求。你的思路方向完全没问题——用共享状态+同步原语来协调各方,下面是业界常用的最佳实践方案:

核心设计:共享同步上下文结构体

首先要定义一个所有Generator和Connection都能访问的共享同步结构体,用来承载协调所需的状态和同步原语:

  • 读写锁:保证dataChunk的写排他(Generator更新时没人能读)、读共享(多个Connection可以同时读)。
  • 活跃连接计数器:原子类型,动态跟踪当前在线的Connection数量,用来确定每轮需要等待多少个消费者完成。
  • 待处理计数器:记录当前还有多少个Connection没完成本轮数据的处理,Generator靠它判断是否可以生成新数据。
  • 两个条件变量:一个给Generator等待所有消费者完成,另一个给Connection等待新数据就绪。
  • 数据版本号:原子类型,用来区分不同轮次的dataChunk,避免消费者重复处理旧数据或者错过新数据。

用C++伪代码示例的话,结构体大概是这样:

struct SyncContext {
    std::shared_mutex data_rw_lock; // 读写锁,保护dataChunk的访问
    std::atomic<int> active_connections; // 当前活跃的Connection数量
    std::atomic<int> pending_processes; // 未完成本轮处理的Connection数
    std::condition_variable_any generator_cv; // Generator等待的条件变量
    std::condition_variable_any connection_cv; // Connection等待的条件变量
    std::atomic<int> data_version; // 数据版本号,每生成一次新数据递增1
    uint8_t* current_data_chunk; // 指向Generator最新生成的dataChunk
};

Generator的工作流程

Generator的核心逻辑是生成数据前等所有消费者处理完,生成后通知所有消费者:

  1. 等待上一轮处理完成:通过条件变量等待pending_processes归0,确保所有Connection都处理完了上一批数据。
  2. 独占更新数据:获取写锁,生成新的dataChunk,更新同步上下文里的current_data_chunk,递增版本号,然后把pending_processes设置为当前的活跃连接数。
  3. 通知消费者:释放写锁,触发Connection的条件变量,告诉所有Connection新数据已经就绪。
  4. 回到步骤1循环。

伪代码示例:

void Generator::run_production_loop() {
    while (true) {
        // 等待所有Connection处理完上一轮数据
        std::unique_lock<std::shared_mutex> lock(sync_ctx.data_rw_lock);
        sync_ctx.generator_cv.wait(lock, [this]() {
            return sync_ctx.pending_processes == 0;
        });

        // 生成新数据(这里是你的generateData逻辑)
        this->generateData();
        sync_ctx.current_data_chunk = this->dataChunk;
        sync_ctx.data_version++;
        // 设置本轮需要等待的Connection数量
        sync_ctx.pending_processes = sync_ctx.active_connections.load();

        // 通知所有Connection新数据已就绪
        sync_ctx.connection_cv.notify_all();
        lock.unlock();
    }
}

Connection的工作流程

Connection需要暂存客户端请求、等待新数据、处理请求、通知Generator完成:

  1. 上线登记:启动时原子递增active_connections,告诉Generator自己加入了消费者队列。
  2. 接收并暂存请求:把收到的客户端请求放到线程安全的队列里,避免等待新数据时丢失请求。
  3. 等待新数据就绪:通过条件变量等待,直到同步上下文的版本号比自己上次处理的版本高(确保是新数据)。
  4. 批量处理请求:获取读锁,读取current_data_chunk的对应部分,处理队列里的所有请求并返回给客户端。
  5. 标记完成并通知:原子递减pending_processes,如果自己是最后一个完成的Connection,就通知Generator可以生成下一批数据。
  6. 循环等待新请求:回到步骤2,继续接收客户端请求。
  7. 下线登记:断开连接时原子递减active_connections。

伪代码示例:

void Connection::run_handler_loop() {
    // 上线:增加活跃连接数
    sync_ctx.active_connections++;
    int last_processed_version = -1;
    std::queue<ClientRequest> request_queue;
    // 这里可以用线程安全队列,比如加锁的std::queue或者无锁队列
    std::mutex queue_mutex;

    while (true) {
        // 接收客户端请求,放入队列
        ClientRequest req = this->receive_client_request();
        std::lock_guard<std::mutex> q_lock(queue_mutex);
        request_queue.push(req);
        q_lock.unlock();

        // 等待新数据就绪(当前版本未处理过)
        std::shared_lock<std::shared_mutex> data_lock(sync_ctx.data_rw_lock);
        sync_ctx.connection_cv.wait(data_lock, [this, &last_processed_version]() {
            return sync_ctx.data_version > last_processed_version;
        });

        // 处理所有暂存的请求
        q_lock.lock();
        while (!request_queue.empty()) {
            ClientRequest current_req = request_queue.front();
            request_queue.pop();

            // 读取dataChunk的指定部分并返回给客户端
            uint8_t* target_data = sync_ctx.current_data_chunk + current_req.offset;
            this->send_data_to_client(target_data, current_req.length);
        }
        q_lock.unlock();

        // 标记当前版本已处理,递减待处理计数器
        last_processed_version = sync_ctx.data_version.load();
        int remaining_processes = --sync_ctx.pending_processes;
        data_lock.unlock();

        // 如果是最后一个完成的,通知Generator可以生成新数据
        if (remaining_processes == 0) {
            sync_ctx.generator_cv.notify_one();
        }
    }

    // 下线:减少活跃连接数
    sync_ctx.active_connections--;
}

关键注意事项

  • 读写锁的选择:如果是C用std::shared_mutex(C17及以上),Java用ReentrantReadWriteLock,其他语言找对应的读写锁实现,核心是保证多读单写。
  • 避免虚假唤醒:条件变量的wait必须带谓词判断(比如上面的sync_ctx.data_version > last_processed_version),不然可能会被虚假唤醒导致错误。
  • 线程安全的请求队列:必须保证多个线程(如果Connection的请求接收是单独线程的话)访问队列时的安全性,要么加锁,要么用无锁队列。
  • 动态连接的处理:Connection上线和下线时的active_connections更新必须是原子操作,避免Generator拿到错误的连接数量。
  • 版本号的必要性:防止Connection在等待期间错过多轮数据,或者重复处理旧数据,确保每轮数据只被处理一次。

这个方案完美覆盖了你的所有需求:

  • Generator只有在所有Connection处理完上一轮后才会生成新数据;
  • 读写锁保证了Connection处理请求时Generator不能写入,同时多个Connection可以并发读取;
  • 请求队列确保了等待新数据时不会丢失客户端请求;
  • 同步上下文的状态让所有对象都能感知系统的运行状态,协调工作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:35:33