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

6线程访问共享队列:线程终止异常与消费线程实现求助

问题分析与解决方案

核心问题定位

1. 前N个线程无法正常终止的原因

  • terminateThreads是普通bool,不具备原子性,线程可能无法及时感知主线程的终止信号,导致循环无法退出
  • 主线程未调用join()等待所有子线程结束,程序退出时触发未定义行为,表现为卡住
  • 原代码中std::max_element(durations, durations + N + 1)越界,durations数组仅包含N个元素,导致sleep时间错误
  • 线程循环中sleep_for在终止条件判断之后,即使收到终止信号,也要等待sleep结束才能退出

2. 第6个线程未实现正确逻辑

  • 原代码复用生产者线程函数,完全未实现消费者监控逻辑
  • 缺少队列等待机制,忙等会浪费CPU资源且可能阻塞生产者
  • 未实现基于计数器的终止条件

修正后的完整代码

#include <iostream>
#include <thread>
#include <queue>
#include <chrono>
#include <atomic>
#include <condition_variable>
#include <algorithm>
#include <string>
#include <random>

const int N = 5;
const int IDs[N] = {0, 1, 2, 3, 4};
const int startDelays[N] = {0, 2, 3, 4, 1};
const int durations[N] = {1, 3, 2, 1, 3};
const int timePeriods[N] = {100, 110, 220, 150, 250};
const int percentageError = 10; // 时间周期误差百分比
const int CONSUMER_TERMINATE_COUNT = 1000; // 消费者终止的消息计数阈值

struct ThreadParams {
    uint8_t id;
    int startDelay;
    int duration;
    int timePeriod;
    char replicationChar;
};

std::queue<std::string> messageQueue;
std::mutex queueMutex;
std::condition_variable queueCV; // 用于消费者等待队列消息
std::atomic<bool> terminateProducers(false); // 原子变量,保证线程间可见性
std::atomic<int> processedCount(0); // 消费者处理的消息计数

// 生产者线程函数(前N个线程)
void producerThread(const ThreadParams& params) {
    std::this_thread::sleep_for(std::chrono::seconds(params.startDelay));
    std::cout << "生产者线程 " << static_cast<char>('A' + params.id) << " 启动." << std::endl;

    int messageNumber = 0;
    auto startTime = std::chrono::system_clock::now();
    std::random_device rd;
    std::mt19937 gen(rd());

    while (!terminateProducers && 
           std::chrono::duration_cast<std::chrono::seconds>(std::chrono::system_clock::now() - startTime).count() < params.duration) {
        if (terminateProducers) break;
        
        // 生成带误差的时间周期
        int errorRange = params.timePeriod * percentageError / 100;
        std::uniform_int_distribution<> dist(params.timePeriod - errorRange, params.timePeriod + errorRange);
        std::this_thread::sleep_for(std::chrono::milliseconds(dist(gen)));

        if (terminateProducers) break;

        // 生成6字节格式的消息:[ID字符, 递增编号, 补全重复字符]
        std::string message;
        message += static_cast<char>('0' + params.id); // ID转为字符(0->'0',1->'1')
        message += std::to_string(messageNumber); // 递增编号,支持多位数
        // 补全到6字节长度
        int fillLen = 6 - message.size();
        if (fillLen > 0) {
            message.append(fillLen, params.replicationChar);
        }

        {
            std::lock_guard<std::mutex> lock(queueMutex);
            messageQueue.push(message);
        }
        queueCV.notify_one(); // 通知消费者有新消息

        std::cout << "生产者线程 " << static_cast<char>('A' + params.id) << " 推送消息: " << message << std::endl;
        ++messageNumber;
    }

    // 推送终止标记消息
    if (!terminateProducers) {
        std::string xMessage;
        xMessage += static_cast<char>('0' + params.id);
        xMessage.append(5, 'X');
        {
            std::lock_guard<std::mutex> lock(queueMutex);
            messageQueue.push(xMessage);
        }
        queueCV.notify_one();
        std::cout << "生产者线程 " << static_cast<char>('A' + params.id) << " 推送终止消息: " << xMessage << std::endl;
    }

    std::cout << "生产者线程 " << static_cast<char>('A' + params.id) << " 退出." << std::endl;
}

// 消费者线程函数(第6个线程)
void consumerThread() {
    std::cout << "消费者线程启动." << std::endl;

    while (processedCount.load() < CONSUMER_TERMINATE_COUNT) {
        std::unique_lock<std::mutex> lock(queueMutex);
        // 等待队列有消息,或生产者全部终止且队列空
        queueCV.wait(lock, []{ 
            return !messageQueue.empty() || terminateProducers; 
        });

        // 生产者终止且队列空时提前退出
        if (terminateProducers && messageQueue.empty()) {
            break;
        }

        // 取出消息后立即解锁,减少锁持有时间
        std::string msg = messageQueue.front();
        messageQueue.pop();
        lock.unlock();

        processedCount++;
        std::cout << "消费者处理消息: " << msg << " | 已处理总数: " << processedCount.load() << std::endl;
    }

    std::cout << "消费者线程达到终止条件,退出." << std::endl;
}

int main() {
    std::thread threads[N + 1];
    ThreadParams producerParams[N];

    // 初始化前N个生产者线程参数
    for (int i = 0; i < N; ++i) {
        producerParams[i].id = static_cast<uint8_t>(i);
        producerParams[i].startDelay = startDelays[i];
        producerParams[i].duration = durations[i];
        producerParams[i].timePeriod = timePeriods[i];
        producerParams[i].replicationChar = static_cast<char>('A' + i);
        threads[i] = std::thread(producerThread, producerParams[i]);
    }

    // 初始化第6个消费者线程
    threads[N] = std::thread(consumerThread);

    // 等待所有生产者线程的最长运行时长
    int maxDuration = *std::max_element(durations, durations + N);
    std::this_thread::sleep_for(std::chrono::seconds(maxDuration + 1)); // 加1秒确保生产者完全结束

    // 通知生产者终止
    terminateProducers = true;
    queueCV.notify_one(); // 唤醒可能等待的消费者

    // 等待所有线程执行完毕
    for (int i = 0; i <= N; ++i) {
        threads[i].join();
    }

    // 处理队列中剩余未处理的消息
    std::cout << "\n队列中剩余未处理消息:" << std::endl;
    std::lock_guard<std::mutex> lock(queueMutex);
    while (!messageQueue.empty()) {
        std::cout << messageQueue.front() << std::endl;
        messageQueue.pop();
    }

    return 0;
}

关键修正点说明

1. 线程终止问题解决

  • 使用std::atomic<bool>替代普通bool,确保终止信号对所有线程可见
  • 生产者线程每次sleep前检查终止信号,避免不必要的阻塞
  • 主线程调用join()等待所有子线程结束,消除未定义行为
  • 修正std::max_element的范围,仅取前N个生产者的运行时长

2. 消费者线程实现

  • 用std::condition_variable实现队列等待机制,避免忙等,节省CPU资源
  • 终止条件为处理消息数达到阈值,或所有生产者终止且队列空
  • 取出消息后立即解锁,减少锁持有时间,保证生产者正常推送

3. 其他优化

  • 修正消息生成逻辑,支持多位数递增编号,确保消息长度为6字节
  • 加入时间周期随机误差,符合原需求
  • 规范线程输出,区分生产者和消费者操作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 16:35:53