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(); }
问题根源
- 线程调度的公平性限制:操作系统默认采用公平线程调度策略,先创建的线程会优先获得CPU时间片。你每次调用
addPersons时会一次性创建批量线程,这些线程会抢占大部分调度资源,新添加的线程很难插队执行。 - 无优先级任务管理:现有代码没有为后添加的人员设置执行优先级,所有线程平等竞争CPU和汽车资源,无法实现新任务优先执行的需求。
- 全局变量数据竞争:
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; }
代码说明
- 优先级任务队列:使用
std::priority_queue实现大顶堆结构,后添加的任务会被赋予更高的优先级数值,确保被优先取出执行。 - 固定工作线程:创建与汽车数量一致的工作线程,避免一次性创建大量线程导致的调度拥堵,线程持续从队列中取任务执行。
- 线程安全同步:用互斥锁和条件变量保护任务队列和输入变量,避免数据竞争和线程阻塞。
- 任务添加逻辑:用户输入后直接将人员任务加入优先级队列,无需创建批量线程,新任务会立即被工作线程抢占执行。
这样修改后,新添加的人员任务会优先获得执行机会,实现你想要的混合分配效果。
内容的提问来源于stack exchange,提问作者Martin Stone
相关产品推荐
相关产品推荐

