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

如何在无限循环中高效并行运行函数(避免重复创建线程)

复用线程实现固定帧率下的并行任务循环

问题背景

需要实现一个固定帧率的无限循环,循环内并行执行多个只读无竞态的函数,原实现每次循环创建/销毁线程导致效率低下,希望复用常驻线程,通过同步机制等待所有任务完成后进入下一次循环。

原有低效实现

#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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 17:15:27