如何在无限循环中高效并行运行函数(避免重复创建线程)
复用线程实现固定帧率下的并行任务循环
问题背景
需要实现一个固定帧率的无限循环,循环内并行执行多个只读无竞态的函数,原实现每次循环创建/销毁线程导致效率低下,希望复用常驻线程,通过同步机制等待所有任务完成后进入下一次循环。
原有低效实现
#include <iostream> #include <thread> #include <random> #include <condition_variable> #include <mutex> int getRandomIntBetween(int minValue, int maxValue) { std::random_device rd; std::mt19937 rng(rd()); std::uniform_int_distribution<int> uni(minValue, maxValue); return uni(rng); } void fun1() { int randomInterval = getRandomIntBetween(10, 90); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun1 done in " << randomInterval << "ms" << std::endl; } void fun2() { int randomInterval = getRandomIntBetween(10, 90); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun2 done in " << randomInterval << "ms" << std::endl; } void fun3() { int randomInterval = getRandomIntBetween(10, 200); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun3 done in " << randomInterval << "ms" << std::endl; } void fun4() { int randomInterval = getRandomIntBetween(3, 300); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun4 done in " << randomInterval << "ms" << std::endl; } int main(int argc, char* argv[]) { const int64_t frameDurationInUs = 1.0e6 / 1; std::cout << "Parallel looping testing" << std::endl; std::condition_variable cv; std::mutex mut; bool stop = false; size_t counter{ 0 }; using delta = std::chrono::duration<int64_t, std::ratio<1, 1000000>>; auto next = std::chrono::steady_clock::now() + delta{ frameDurationInUs }; std::unique_lock<std::mutex> lk(mut); while (!stop) { mut.unlock(); if (counter % 10 == 0) { std::cout << counter << " frames..." << std::endl; } std::thread t1{ &fun1 }; std::thread t2{ &fun2 }; std::thread t3{ &fun3 }; std::thread t4{ &fun4 }; counter++; t1.join(); t2.join(); t3.join(); t4.join(); mut.lock(); cv.wait_until(lk, next); next += delta{ frameDurationInUs }; } return 0; }
优化实现:复用常驻线程
核心思路是让每个线程保持活跃,通过条件变量+状态标记控制任务触发与同步:
- 每个线程循环等待主线程的"开始任务"信号
- 主线程触发所有线程执行任务后,等待所有线程完成
- 完成后进入下一次帧率控制循环
#include <iostream> #include <thread> #include <random> #include <condition_variable> #include <mutex> #include <vector> #include <functional> #include <algorithm> // 线程状态枚举 enum class ThreadState { Idle, // 等待任务触发 Running, // 正在执行任务 Done // 任务执行完成 }; int getRandomIntBetween(int minValue, int maxValue) { static thread_local std::random_device rd; static thread_local std::mt19937 rng(rd()); std::uniform_int_distribution<int> uni(minValue, maxValue); return uni(rng); } void fun1() { int randomInterval = getRandomIntBetween(10, 90); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun1 done in " << randomInterval << "ms" << std::endl; } void fun2() { int randomInterval = getRandomIntBetween(10, 90); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun2 done in " << randomInterval << "ms" << std::endl; } void fun3() { int randomInterval = getRandomIntBetween(10, 200); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun3 done in " << randomInterval << "ms" << std::endl; } void fun4() { int randomInterval = getRandomIntBetween(3, 300); std::this_thread::sleep_for(std::chrono::milliseconds(randomInterval)); std::cout << "fun4 done in " << randomInterval << "ms" << std::endl; } // 工作线程函数:接收任务、状态标记、同步锁和条件变量 void workerThread(std::function<void()> task, ThreadState& state, std::mutex& mutex, std::condition_variable& cv, bool& stop) { while (!stop) { std::unique_lock<std::mutex> lk(mutex); // 等待主线程触发任务,或者停止信号 cv.wait(lk, [&](){ return state == ThreadState::Running || stop; }); if (stop) break; lk.unlock(); // 执行任务 task(); lk.lock(); state = ThreadState::Done; // 通知主线程任务完成 cv.notify_all(); } } int main(int argc, char* argv[]) { const int64_t frameDurationInUs = 1.0e6 / 1; std::cout << "Parallel looping testing" << std::endl; bool stop = false; size_t counter = 0; using delta = std::chrono::duration<int64_t, std::ratio<1, 1000000>>; auto next = std::chrono::steady_clock::now() + delta{frameDurationInUs}; // 初始化线程状态、锁和条件变量 std::mutex mutex; std::condition_variable cv; std::vector<ThreadState> states = {ThreadState::Idle, ThreadState::Idle, ThreadState::Idle, ThreadState::Idle}; std::vector<std::function<void()>> tasks = {fun1, fun2, fun3, fun4}; std::vector<std::thread> threads; // 创建常驻工作线程 for (size_t i = 0; i < tasks.size(); ++i) { threads.emplace_back(workerThread, tasks[i], std::ref(states[i]), std::ref(mutex), std::ref(cv), std::ref(stop)); } std::unique_lock<std::mutex> lk(mutex); while (!stop) { lk.unlock(); if (counter % 10 == 0) { std::cout << counter << " frames..." << std::endl; } lk.lock(); // 设置所有线程为运行状态,触发任务执行 for (auto& state : states) { state = ThreadState::Running; } cv.notify_all(); // 等待所有线程完成任务 cv.wait(lk, [&](){ return std::all_of(states.begin(), states.end(), [](const ThreadState& s){ return s == ThreadState::Done; }); }); // 重置所有线程状态为空闲,准备下一轮 for (auto& state : states) { state = ThreadState::Idle; } counter++; lk.unlock(); // 固定帧率控制 std::this_thread::sleep_until(next); next += delta{frameDurationInUs}; } // 停止所有线程 lk.lock(); stop = true; cv.notify_all(); lk.unlock(); // 等待所有线程退出 for (auto& t : threads) { t.join(); } return 0; }
关键优化点说明
- 常驻线程复用:仅在程序启动时创建一次线程,避免反复创建销毁的开销
- 状态同步:用
ThreadState标记每个线程的状态,主线程通过条件变量触发任务并等待所有线程完成 - 线程局部随机数:将
std::random_device和std::mt19937改为thread_local,避免多线程竞争随机数生成器,提升效率 - 安全退出:通过
stop标记优雅终止所有工作线程
内容的提问来源于stack exchange,提问作者Jepessen
相关产品推荐
相关产品推荐

