C++ / Robotics · 系统工程 · LESSON 20

消息队列与生产者消费者

把传感器输入、处理任务和输出消息拆成有界队列,并处理背压和退出。

18 分钟message queue · producer consumer · backpressure

消息队列与生产者消费者

Promise 队列可以表达异步任务,但传感器以固定频率产生数据,处理速度一旦下降就会形成背压。无界 JS 数组会把问题变成内存增长;C++ 队列应明确容量、满时动作、消息顺序、关闭状态和所有权。队列不是单纯的数据结构,而是系统在过载时如何降级的策略。

学习目标

你将学会解释 producer/consumer、背压、容量和消息所有权;实现阻塞或非阻塞的 push/pop 与关闭协议;区分最新状态、事件和点云的丢弃策略;把队列延迟接入机器人指标;用超时、Sanitizer 和 TSan 排查退出卡死与并发错误。

容量把延迟变成可观测上限

TRANSLATION LENS 同一个意图,两种工程表达 窄屏可左右滑动查看完整代码
JS / TS
const queue = [];
function onMessage(message) {
queue.push(message);
drain();
}
C++
BoundedQueue<Message> queue{64};
void producer(Message message) { queue.try_push(std::move(message)); }
void consumer() {
while (auto message = queue.pop()) process(std::move(*message));
}

容量 64 不代表一定正确:若消息每 10ms 到达、处理平均 30ms,队列很快会满。状态流可以覆盖旧值或只保留最新值,安全事件可能必须阻塞生产者并报警,点云则可能降低采样或丢弃最旧帧。把满队列计数和最大深度作为诊断指标,才能判断调参是否有效。

#include <optional>

template <typename T>
class LatestValue {
 public:
  void put(T value) { value_ = std::move(value); }
  std::optional<T> take() {
    if (!value_) return std::nullopt;
    auto result = std::move(value_);
    value_.reset();
    return result;
  }
 private:
  std::optional<T> value_;
};

这是单线程的最新状态槽:连续 put(1)put(2)put(3)take() 应得到 3。它适合当前距离或姿态,不适合碰撞事件,因为旧事件会被覆盖。跨线程使用时仍要补 mutex 或放进一个已有同步的 owner 中。

一个可关闭的有界队列

#include <condition_variable>
#include <deque>
#include <mutex>
#include <optional>

template <class T>
class BoundedQueue {
public:
  explicit BoundedQueue(std::size_t capacity) : capacity_{capacity} {}
  bool try_push(T value) {
    std::scoped_lock lock{mutex_};
    if (closed_ || items_.size() >= capacity_) return false;
    items_.push_back(std::move(value));
    ready_.notify_one();
    return true;
  }
  std::optional<T> pop() {
    std::unique_lock lock{mutex_};
    ready_.wait(lock, [&] { return closed_ || !items_.empty(); });
    if (items_.empty()) return std::nullopt;
    T value = std::move(items_.front());
    items_.pop_front();
    return value;
  }
  void close() { std::scoped_lock lock{mutex_}; closed_ = true; ready_.notify_all(); }
private:
  std::size_t capacity_;
  std::deque<T> items_;
  std::mutex mutex_;
  std::condition_variable ready_;
  bool closed_{};
};

pop 在 close 后仍会先排空已有消息;若业务要求立即丢弃,可在 close 中清空并记录数量。生产者和消费者都要把 false/空结果当作协议的一部分,不能在关闭后继续 push,也不能让析构直接销毁仍有线程访问的队列。

关闭、排空与析构顺序

void close() {
  std::scoped_lock lock{mutex_};
  closed_ = true;
  ready_.notify_all();
}

// consumer
while (auto message = queue.pop()) {
  process(std::move(*message));
}

关闭应先停止 producer,再标记队列并唤醒等待者,consumer 排空允许保留的消息后退出,最后 owner 才能析构队列和设备。若要立即丢弃,必须记录 dropped 数量并说明数据为何可以丢。验证时给队列三个带序号消息,close 后确认尾部顺序和线程退出时间。

按值、移动和引用的生命周期

struct Detection { std::uint64_t sequence{}; std::vector<float> points; };

void publish(BoundedQueue<Detection>& queue, Detection message) {
  if (!queue.try_push(std::move(message))) record_drop();
}

按值接收再移动,成功入队后队列拥有 points;失败时调用者可能只剩 moved-from 对象,因此接口应明确失败后的所有权,或在尝试前复制/使用 unique_ptr。把回调参数的引用存进异步队列是常见悬空 bug,消息必须拥有自己的数据或绑定到有明确生命周期的共享存储。

所有权、阻塞和丢弃

按值 push 接收消息后移动入队,队列成为消息的拥有者;引用入队则必须保证引用对象持续存在,通常不适合跨线程。阻塞 push 可以提供背压,却可能让采集线程错过硬件截止时间;try_push 丢弃可保持采集实时性,但必须记录原因。根据消息类型分别选择,不要整套系统套同一种队列。

JS/TS 迁移反例:数组不是协议

用共享 std::deque 模仿 JS 数组并不能自动解决并发;push_backpop_front 需要同一把锁,析构前还要停止生产者。也不要无界地积压消息来“保证不丢数据”:点云每帧数 MB 时,延迟会先变大,最后触发内存压力。队列策略必须把“丢什么、何时丢、如何报警”写成 API 或指标。

编译和运行时排错

条件变量缺少谓词会被虚假唤醒;等待期间用 items_.empty() 再检查。消费者退出不了,检查 close 是否 notify_all;队列偶发超容量,检查所有读写是否都在同一 mutex 内。延迟持续升高时看 enqueue/dequeue 时间戳、满队列次数和处理耗时,别只把容量调大。Sanitizer 可查队列对象被提前销毁,TSan 可查漏锁访问。

连接传感器与消息系统

采集端推入带 stamp_nssequenceframe_id 的消息;处理端在 dequeue 时计算 age;发布端记录处理结果和失败原因。距离状态可以容量 1 丢旧值,碰撞事件用 FIFO 并在满时进入告警,点云可在背压时丢最旧帧但保留统计。回放测试要验证顺序、最大年龄、dropped_count 和 close 后的尾部消息。

迁移练习

为距离状态建立容量 1 的 latest-value 槽,为碰撞事件建立容量 64 的 FIFO。写出生产者、消费者和 close 顺序,分别决定满时覆盖、拒绝还是阻塞,并为每种结果增加 dropped/blocked 计数。

01
TRY IT YOURSELF

设计两种消息队列策略

使用或改写 BoundedQueue,说明距离、点云和碰撞事件在容量、顺序、满队列和关闭上的差异。

给我一点提示

控制状态追求新鲜度,事件追求完整性;阻塞采集线程前要评估实时截止时间。

查看参考答案
距离可用 latest-value,只保留最新样本并记录覆盖次数;碰撞事件用有界 FIFO,满时拒绝并报警或进入安全状态;点云可限深并丢弃最旧帧。close 唤醒所有等待者,消费者处理完允许保留的消息后退出。
本节结论

一个队列的正确性包含数据结构和故障策略两部分。下一节会把容量、布局与缓存访问联系起来,测量“快”到底来自哪里。

FURTHER READING

延伸阅读

先完成本节练习,再用这些资料查阅完整 API 和真实项目组织方式。

当前学习阶段系统工程
0/7

阶段共 7 节课,按顺序完成更容易建立完整的迁移模型。