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

满足条件时并发执行多组任务的设计模式及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完成
  } 
}

咨询点:

  1. 是否有对应设计模式解决这类「分组线程满足特定条件才并发执行,否则等待」的问题?
  2. C++有哪些特性可实现该行为?
  3. 当前设计是否存在缺陷?
  4. 如何用更优雅的方式实现需求?

对应的设计模式

  • 栅栏(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等函数,任务调度逻辑和业务逻辑耦合在一起,不利于维护。

更优雅的实现方案

  1. 封装任务调度器:创建一个ControlSystem类,内部维护执行器列表、传感器列表和同步原语(比如std::barrier),把调度逻辑封装在类里,主线程只需要调用ControlSystem::runLoop()即可。
  2. 复用线程池:提前创建固定数量的线程,循环从任务队列中取执行器/传感器的step()任务执行,避免每次循环创建线程的开销。可以自己基于std::thread和队列实现,也用成熟的线程池逻辑。
  3. 异常安全:在并发执行任务时,捕获step()抛出的异常,记录日志并标记故障组件,避免影响整个系统运行。
  4. 扩展友好:用容器(比如std::vector<std::unique_ptr<GenericActuator>>)管理执行器和传感器,新增组件时只需要向容器中添加实例,无需修改调度逻辑。
  5. 状态监控:为执行器/传感器添加状态接口(比如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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 23:09:55