满足条件时并发执行多组任务的设计模式及C++实现问询
问题描述
模拟一个控制系统,核心是无限循环执行三组任务,规则如下:
- 执行器(actuators):多个且完全解耦,可并行执行,必须全部完成后才能进入下一阶段
- 动力学模块(dynamics):单任务,需在所有执行器完成后运行
- 传感器(sensors):多个且完全解耦,可并行执行,需在动力学模块完成后运行
伪C++代码示例:
class GenericActuator { virtual void step() = 0; }; class Actuator1 : GenericActuator { void step() override; // other methods }; class Actuator2 : GenericActuator { void step() override; // other methods }; class GenericSensor { virtual void step() = 0; }; class Sensor1 : GenericSensor { void step() override; // other methods }; class Sensor2 : GenericSensor { void step() override; // other methods }; int main() { while (true) { runConcurrentlyStepForActuators(); // 阻塞直到所有执行器step完成 runDynamics(); // 单函数,无并发问题 runConcurrentlyStepForSensors(); // 阻塞直到所有传感器step完成 } }
咨询点:
- 是否有对应设计模式解决这类「分组线程满足特定条件才并发执行,否则等待」的问题?
- C++有哪些特性可实现该行为?
- 当前设计是否存在缺陷?
- 如何用更优雅的方式实现需求?
对应的设计模式
- 栅栏(Barrier)模式:完全匹配场景——一组线程必须全部到达某个同步点后,才能继续执行后续逻辑。这里每个阶段(执行器、传感器)的并发任务都需要等待全部完成,才能进入下一个阶段,栅栏正好用来做这种同步。
- 流水线(Pipeline)模式:整个控制系统是一个循环的流水线,分为「执行器并行处理」→「动力学计算」→「传感器并行处理」三个阶段,每个阶段必须在上一阶段完全结束后才启动,符合流水线的阶段依赖特性。
C++特性实现方式
- C++20
std::barrier:可重复使用的同步原语,完美适配无限循环场景。每个执行器线程在step()完成后调用barrier.arrive_and_wait(),所有线程都到达后,栅栏自动重置,下一轮循环可以继续使用。 - C++11+
std::async+std::future:在runConcurrentlyStepForActuators里,为每个执行器创建异步任务(std::async(std::launch::async, &GenericActuator::step, actuator)),然后把所有std::future存入容器,遍历调用future.get()等待全部完成。这种方式不用手动管理线程,适合简单场景,但每次循环创建线程可能有开销。 - C++11+
std::thread+std::condition_variable+std::mutex:手动实现栅栏逻辑——维护一个计数器,每个任务完成后计数器减一,当计数器归零时,通知主线程继续执行。适合C++20之前的版本。 - C++20
std::latch:如果是一次性同步可以用,但因为是无限循环,每次循环都要重新创建latch,不如std::barrier高效。
当前设计的缺陷
- 扩展性差:
runConcurrentlyStepForActuators和runConcurrentlyStepForSensors是硬编码的黑盒,新增执行器/传感器时必须修改这两个函数的内部逻辑,违反开闭原则。 - 缺乏异常处理:如果某个执行器/传感器的
step()抛出异常,当前设计没有处理逻辑,可能导致整个系统崩溃或进入未知状态。 - 线程资源浪费:如果每次循环都创建新线程执行任务,频繁的线程创建销毁会带来额外开销,降低性能。
- 无状态管理:没有考虑执行器/传感器的就绪状态,比如某个执行器故障无法执行
step(),系统无法感知并做出处理。 - 耦合性高:主线程直接调用
runConcurrentlyStepForActuators等函数,任务调度逻辑和业务逻辑耦合在一起,不利于维护。
更优雅的实现方案
- 封装任务调度器:创建一个
ControlSystem类,内部维护执行器列表、传感器列表和同步原语(比如std::barrier),把调度逻辑封装在类里,主线程只需要调用ControlSystem::runLoop()即可。 - 复用线程池:提前创建固定数量的线程,循环从任务队列中取执行器/传感器的
step()任务执行,避免每次循环创建线程的开销。可以自己基于std::thread和队列实现,也用成熟的线程池逻辑。 - 异常安全:在并发执行任务时,捕获
step()抛出的异常,记录日志并标记故障组件,避免影响整个系统运行。 - 扩展友好:用容器(比如
std::vector<std::unique_ptr<GenericActuator>>)管理执行器和传感器,新增组件时只需要向容器中添加实例,无需修改调度逻辑。 - 状态监控:为执行器/传感器添加状态接口(比如
isReady()),调度前检查状态,跳过故障组件或触发告警。
示例简化实现:
#include <vector> #include <memory> #include <barrier> #include <thread> #include <stdexcept> #include <iostream> class GenericActuator { public: virtual void step() = 0; virtual bool isReady() const = 0; virtual ~GenericActuator() = default; }; class Actuator1 : public GenericActuator { public: void step() override { /* 执行具体逻辑 */ } bool isReady() const override { return true; } }; class GenericSensor { public: virtual void step() = 0; virtual bool isReady() const = 0; virtual ~GenericSensor() = default; }; class Sensor1 : public GenericSensor { public: void step() override { /* 执行具体逻辑 */ } bool isReady() const override { return true; } }; class ControlSystem { public: void addActuator(std::unique_ptr<GenericActuator> actuator) { actuators_.push_back(std::move(actuator)); // 更新栅栏的参与线程数 actuator_barrier_ = std::barrier(actuators_.size()); } void addSensor(std::unique_ptr<GenericSensor> sensor) { sensors_.push_back(std::move(sensor)); sensor_barrier_ = std::barrier(sensors_.size()); } void runLoop() { while (true) { runActuators(); runDynamics(); runSensors(); } } private: void runActuators() { std::vector<std::thread> threads; for (auto& act : actuators_) { if (!act->isReady()) continue; threads.emplace_back([this, &act]() { try { act->step(); actuator_barrier_.arrive_and_wait(); } catch (const std::exception& e) { std::cerr << "Actuator error: " << e.what() << std::endl; // 即使异常,也要通知栅栏,避免其他线程永久等待 actuator_barrier_.arrive_and_drop(); } }); } for (auto& t : threads) t.join(); } void runDynamics() { /* 动力学模块具体逻辑 */ } void runSensors() { std::vector<std::thread> threads; for (auto& sen : sensors_) { if (!sen->isReady()) continue; threads.emplace_back([this, &sen]() { try { sen->step(); sensor_barrier_.arrive_and_wait(); } catch (const std::exception& e) { std::cerr << "Sensor error: " << e.what() << std::endl; sensor_barrier_.arrive_and_drop(); } }); } for (auto& t : threads) t.join(); } std::vector<std::unique_ptr<GenericActuator>> actuators_; std::vector<std::unique_ptr<GenericSensor>> sensors_; std::barrier<> actuator_barrier_{0}; std::barrier<> sensor_barrier_{0}; }; int main() { ControlSystem cs; cs.addActuator(std::make_unique<Actuator1>()); cs.addSensor(std::make_unique<Sensor1>()); // 添加其他执行器、传感器... cs.runLoop(); }
内容的提问来源于stack exchange,提问作者lumaca96
相关产品推荐
相关产品推荐

