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

C++多线程调度问题:如何让后添加线程与先添加线程混合执行

问题描述

我用C++写了一个多线程程序,核心逻辑是把人员分配给汽车——每辆汽车是被线程(代表人员)锁定的资源。程序支持运行时添加人员线程,这些线程会抢占汽车资源。

现在遇到的问题是:先添加50个人员线程,再添加20个时,这20个线程必须等50个全部完成分配才会开始执行。实际执行顺序如下:

50 <- 添加50个人员
来自组50的人员被分配到汽车0
来自组50的人员被分配到汽车1
20 <- 添加20个人员
来自组50的人员被分配到汽车2
……(组50的所有人员完成分配)
来自组20的人员被分配到汽车0
来自组20的人员被分配到汽车1
……

我想要实现混合执行的效果:无论何时添加人员线程,后添加的线程都能插队到之前的线程前面执行,比如:

50 <- 添加50个人员
来自组50的人员被分配到汽车0
来自组50的人员被分配到汽车1
20 <- 添加20个人员
来自组50的人员被分配到汽车2
来自组20的人员被分配到汽车3
来自组20的人员被分配到汽车4
来自组50的人员被分配到汽车0
来自组20的人员被分配到汽车1
……

我试过用detach()但没达到预期,相关代码如下:

int numberOfCars = 5;
int numberOfPersons;
bool typed = false;

void addPersons(int numberOfPersons)
{
    std::vector<Person> persons = makePersons(numberOfPersons);
    std::vector<std::thread> threads;
    for (int i = 0; i < persons.size(); i++) {
        int carId = i % numberOfCars;
        threads.emplace_back(
            [carId, person = std::move(persons[i])]() mutable { // 修正原代码的person[i]拼写错误
                car[carId].addPerson(person);
            });
    }
    for (int i = 0; i < threads.size(); i++) {
        threads[i].join();
    }
}

void userInput()
{
    while (true) {
        std::string inputVal;
        std::cin >> inputVal;
        numberOfPersons = std::stoi(inputVal); // 修正原代码的类型转换错误
        typed = true;
    }
}

int main()
{
    std::thread inputThread(userInput);
    std::vector<Car> cars(numberOfCars); // 补充原代码缺失的cars定义
    for (int i = 0; i < numberOfCars; i++) // 修正原代码的冒号语法错误
        cars[i].setId(i + 1);

    while (true) {
        if (typed) {
            typed = false;
            std::thread(addPersons, numberOfPersons).detach(); // 修正原代码的addPerson拼写错误
        }
    }

    inputThread.join();
    return 0;
}

补充的Car类锁相关函数(修正语法错误后):

void Car::lockCar() {
    std::unique_lock<std::mutex> lock(carMutex);
    carFreeCV.wait(lock, [this]() { return isFree; });
    isFree = false;
}

void Car::unlockCar()
{ 
    {
        std::lock_guard<std::mutex> lock(carMutex);
        isFree = true;
    }
    carFreeCV.notify_one();
}

void Car::addPerson(Person& person) {
    lockCar();
    std::cout << "Person: " << person.name() << " added to car: " << id << std::endl;
    std::this_thread::sleep_for(std::chrono::seconds(2)); // 修正原代码的括号缺失错误
    unlockCar();
}
问题根源
  1. 线程调度的公平性限制:操作系统默认采用公平线程调度策略,先创建的线程会优先获得CPU时间片。你每次调用addPersons时会一次性创建批量线程,这些线程会抢占大部分调度资源,新添加的线程很难插队执行。
  2. 无优先级任务管理:现有代码没有为后添加的人员设置执行优先级,所有线程平等竞争CPU和汽车资源,无法实现新任务优先执行的需求。
  3. 全局变量数据竞争:numberOfPersons和typed是未加保护的全局变量,userInput线程写入、main线程读取时存在数据竞争,可能导致线程创建时机异常。
解决方案:基于优先级任务队列的调度

要实现新添加人员优先执行的效果,我们需要用优先级任务队列管理所有待分配的人员,让后添加的任务排在队列头部,同时用固定数量的工作线程消费队列任务,确保新任务被优先处理。

修正后的完整代码示例

#include <iostream>
#include <vector>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <queue>
#include <functional>
#include <chrono>
#include <string>

// Person类定义
class Person {
private:
    std::string name_;
public:
    Person(std::string name) : name_(std::move(name)) {}
    const std::string& name() const { return name_; }
};

// Car类
class Car {
private:
    int id_;
    bool isFree_ = true;
    std::mutex carMutex_;
    std::condition_variable carFreeCV_;
public:
    void setId(int id) { id_ = id; }
    int id() const { return id_; }

    void lockCar() {
        std::unique_lock<std::mutex> lock(carMutex_);
        carFreeCV_.wait(lock, [this]() { return isFree_; });
        isFree_ = false;
    }

    void unlockCar() {
        {
            std::lock_guard<std::mutex> lock(carMutex_);
            isFree_ = true;
        }
        carFreeCV_.notify_one();
    }

    void addPerson(const Person& person) {
        lockCar();
        std::cout << "Person: " << person.name() << " added to car: " << id_ << std::endl;
        std::this_thread::sleep_for(std::chrono::seconds(2));
        unlockCar();
    }
};

// 全局配置与同步
const int numberOfCars = 5;
std::vector<Car> cars(numberOfCars);
std::mutex inputMutex_;

// 任务结构:包含人员信息和优先级(数值越大优先级越高)
struct Task {
    Person person;
    int priority;

    // 重载运算符,让优先级队列优先取出高优先级任务
    bool operator<(const Task& other) const {
        return priority < other.priority;
    }
};

// 优先级任务队列与同步
std::priority_queue<Task> taskQueue_;
std::mutex queueMutex_;
std::condition_variable queueCV_;
bool stopWorkers_ = false;

// 工作线程:负责从队列取任务并分配到汽车
void workerThread(int carId) {
    while (!stopWorkers_) {
        Task task;
        {
            std::unique_lock<std::mutex> lock(queueMutex_);
            queueCV_.wait(lock, []() { return stopWorkers_ || !taskQueue_.empty(); });
            if (stopWorkers_) break;
            
            // 取出队列头部的高优先级任务
            task = std::move(const_cast<Task&>(taskQueue_.top()));
            taskQueue_.pop();
        }
        cars[carId].addPerson(task.person);
    }
}

// 生成Person实例的工具函数
std::vector<Person> makePersons(int count, const std::string& groupTag) {
    std::vector<Person> persons;
    for (int i = 0; i < count; ++i) {
        persons.emplace_back(groupTag + "_" + std::to_string(i));
    }
    return persons;
}

// 添加人员任务到优先级队列
void addPersons(int count, int priority) {
    auto persons = makePersons(count, "group_" + std::to_string(priority));
    {
        std::lock_guard<std::mutex> lock(queueMutex_);
        for (auto& person : persons) {
            taskQueue_.push({std::move(person), priority});
        }
    }
    queueCV_.notify_all(); // 通知工作线程有新任务
}

// 用户输入处理线程
void userInput() {
    static int priorityCounter = 0; // 每次输入的优先级递增,确保后输入的任务优先级更高
    while (true) {
        std::string inputVal;
        std::cin >> inputVal;
        int count = std::stoi(inputVal);
        
        std::lock_guard<std::mutex> lock(inputMutex_);
        addPersons(count, ++priorityCounter);
    }
}

int main() {
    // 初始化汽车ID
    for (int i = 0; i < numberOfCars; ++i) {
        cars[i].setId(i + 1);
    }

    // 创建与汽车数量匹配的工作线程
    std::vector<std::thread> workers;
    for (int i = 0; i < numberOfCars; ++i) {
        workers.emplace_back(workerThread, i);
    }

    // 启动输入线程
    std::thread inputThread(userInput);
    inputThread.join();

    // 停止工作线程
    {
        std::lock_guard<std::mutex> lock(queueMutex_);
        stopWorkers_ = true;
    }
    queueCV_.notify_all();
    for (auto& worker : workers) {
        worker.join();
    }

    return 0;
}

代码说明

  1. 优先级任务队列:使用std::priority_queue实现大顶堆结构,后添加的任务会被赋予更高的优先级数值,确保被优先取出执行。
  2. 固定工作线程:创建与汽车数量一致的工作线程,避免一次性创建大量线程导致的调度拥堵,线程持续从队列中取任务执行。
  3. 线程安全同步:用互斥锁和条件变量保护任务队列和输入变量,避免数据竞争和线程阻塞。
  4. 任务添加逻辑:用户输入后直接将人员任务加入优先级队列,无需创建批量线程,新任务会立即被工作线程抢占执行。

这样修改后,新添加的人员任务会优先获得执行机会,实现你想要的混合分配效果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 03:44:57