线程与消息循环
从事件循环出发,建立机器人消息处理与线程同步的基础模型。
线程与消息循环
Node.js 开发者习惯一个事件循环和回调队列;C++ 机器人程序可能有采集线程、处理线程和发布线程,因此必须回答更多问题:谁拥有可变状态,消息何时过期,队列是否有上限,关闭如何唤醒等待者。并发设计首先是数据流和时间策略,线程 API 只是实现选择。
学习目标
你将学会创建和回收 std::thread/std::jthread;设计明确的 stop 协议;理解 join、detach 与引用捕获的生命周期关系;把传感器生产者和处理消费者隔开;通过序号、线程 ID、队列深度和耗时验证并发系统,而不是用增加 sleep 掩盖竞态。
把事件循环拆成明确的消费者
queue.on("message", (message) => {
handle(message);
});
controller.on("close", () => queue.stop()); std::jthread worker([&](std::stop_token stop) {
while (!stop.stop_requested()) {
if (auto message = queue.wait_pop(stop)) handle(*message);
}
}); std::jthread 的 RAII 让线程离开作用域时请求停止并等待结束;队列的 wait_pop 还必须响应 stop token 或 close 状态,否则 worker 可能永远睡在条件变量上。不要把“有一个全局 running 布尔值”当作完整协议,关闭顺序要包含停止生产、唤醒消费者和排空/丢弃剩余消息。
#include <chrono>
#include <iostream>
#include <thread>
void poll(std::stop_token stop) {
while (!stop.stop_requested()) {
std::cout << "sample\n";
std::this_thread::sleep_for(std::chrono::milliseconds{10});
}
}
int main() {
std::jthread sensor{poll};
std::this_thread::sleep_for(std::chrono::milliseconds{35});
}
用 g++ -std=c++20 -pthread worker.cpp -o worker 编译,程序退出时应由 jthread 请求停止并 join。这里的 sleep 只模拟设备等待,不能当作真实采样时钟;实际驱动应有可取消的阻塞调用。
一个可停止的工作线程
#include <chrono>
#include <thread>
void process_loop(std::stop_token stop) {
while (!stop.stop_requested()) {
// 从有界队列等待一条消息;等待应有超时或关闭唤醒
std::this_thread::sleep_for(std::chrono::milliseconds{1});
process_one_message();
}
}
int main() {
std::jthread worker(process_loop);
run_until_shutdown();
}
示例中的 sleep 只是展示取消点,不是实时同步方案;生产代码用条件变量、消息队列或中间件等待。工作线程不要直接修改节点所有共享字段,尽量让它消费值消息并独占局部状态,结果再通过队列交回发布线程。
条件变量的正确等待
#include <condition_variable>
#include <mutex>
#include <queue>
std::mutex mutex;
std::condition_variable ready;
std::queue<int> messages;
bool closed = false;
bool wait_pop(int& value) {
std::unique_lock lock{mutex};
ready.wait(lock, [] { return closed || !messages.empty(); });
if (messages.empty()) return false;
value = messages.front(); messages.pop();
return true;
}
条件变量允许虚假唤醒,所以一定要使用谓词并在醒来后重新检查状态。生产者 push 后 notify_one,关闭时设置 closed 并 notify_all。验证时推入序号 1、2、3,确认顺序,再 close 并确认等待线程返回 false;没有通知的关闭会在 join 超时处暴露。
三阶段数据流
struct Detection { std::uint64_t sequence{}; float score{}; };
void process(Detection message) {
if (message.score > 0.8F) publish_alert(message.sequence);
}
采集线程只生成完整消息,处理线程只操作自己的局部状态,发布线程处理网络或 ROS 2 输出。跨线程交接值或拥有指针,不要把设备 buffer 的引用塞进异步队列。这样可以分别观察采集、排队、处理和发布延迟,也能在发布端断开时保持采集端可控。
队列策略先于线程数量
控制当前姿态通常只需最新值;点云离线分析需要按序保留;诊断事件可以有界排队并计数丢弃。处理耗时超过输入周期时,不能无限增大队列来隐藏背压,否则延迟会越来越大。为每类消息写容量、顺序、过期和失败动作,性能问题才有可观测指标。
JS/TS 的迁移反例
Promise/async 不等于 C++ 线程安全;两个 callback 仍可能同时修改普通 vector。不要捕获局部变量引用后启动 detach 线程,也不要让一个全局 bool 代替关闭、唤醒和资源析构协议。每个线程必须有 owner、停止点、唤醒机制和 join 边界;每条消息必须说明由谁拥有和何时过期。
编译和运行时排错
std::jthread 找不到说明标准版本低于 C++20;joinable 或析构挂起说明线程仍在阻塞等待,检查关闭信号是否被传到队列。偶发崩溃可能是数据竞争,ASan 不会替你发现所有竞争,要用 TSan 或审查锁/原子协议。CPU 占满而没有消息通常是忙等;延迟尖峰则检查锁竞争、队列分配和日志 I/O,别先增加 worker 数量。
观测并验证并发行为
每个消息带 sequence、采集时间和 frame_id,每个阶段记录线程 ID、开始/结束时间和队列深度。测试覆盖无消息关闭、生产者先停、消费者排空、处理超时和发布失败。用 TSan 跑固定回放数据,用 ASan 检查线程退出后是否仍访问对象;用重复运行确认顺序和结果符合协议,而不是只看一次成功。
迁移练习
设计一个采集线程每 100ms 产生距离消息、处理线程偶尔耗时 300ms 的系统。分别为控制输出和离线记录选择 latest-value、有限 FIFO 或有界丢弃队列,写出停止顺序,并实现一个可响应关闭的 std::jthread 循环。
为消息循环写出关闭协议
确定生产者停止、队列唤醒、消费者退出和资源析构的顺序;为两种消息说明是否允许丢失、是否保序和如何记录背压。
给我一点提示
控制状态重视最新值,审计数据重视完整性;任何有界策略都要记录丢弃或拒绝计数。
查看参考答案
先停止生产,再关闭或唤醒队列,让 wait_pop 返回关闭状态;worker 处理完允许的数据后退出,jthread join,最后析构设备。控制消息可用只保留最新值并记录覆盖次数;离线消息用有界 FIFO,满时拒绝并记录丢弃,不能无限增长。 本节结论
并发系统的可靠性来自可见的生命周期和时间策略,而不是线程越多越快。下一节会用 mutex、atomic 和内存序把共享状态的规则写得更精确。
延伸阅读
先完成本节练习,再用这些资料查阅完整 API 和真实项目组织方式。
阶段共 7 节课,按顺序完成更容易建立完整的迁移模型。