设计:支持在 Rosbag2 recorder 中重复 transient-local("锁存")消息
作者:Michael Orlov
一些 ROS 2 话题使用 Transient Local 持久性(ROS 1 的”锁存/latched”),使迟到的订阅者也能接收到已发布的数据,但 Rosbag2 用户仍面临两个实际缺口:
- 拆分 bag:消费者通常期望每个拆分文件的开头有”锁存”状态(例如
/tf_static、/map),就像 ROS 1 的--repeat-latched。Rosbag2 目前不提供这个功能。 - 快照模式:当循环缓存覆盖早期的 transient-local 消息时,快照可能丢失它们(例如缓存压力下
/tf_static丢失)。
ROS 2 的”锁存可见性”要求发布者和订阅者都使用 Transient Local 持久性。 当默认值不匹配时,Rosbag2 支持通过 QoS 覆盖文件来处理录制/回放。
- 为 recorder 拆分和快照提供可选的”重复锁存”体验。
- 修复快照模式在缓存压力下对选定的 transient-local 话题的驱逐问题。
- 保持默认行为的确定性:“重复”必须被显式请求。也就是说,transient-local 消息默认不会被重复,否则会引入录制时并未真正发布的额外消息。
- 自动检测所有 transient-local 话题并默认重复它们。
用户可见的接口
Section titled “用户可见的接口”添加 ros2 bag record --repeat-transient-local 并支持按话题指定消息数量,例如:
$ ros2 bag record -a --repeat-transient-local /map=1 /tf_static=5语义:
- 每个条目
<topic>=N表示”捕获/<topic>最近观察到的 N 条消息,并在拆分创建或快照写入创建的每个新 bag 文件开头重新发出”。 - 如果省略
=N,默认 N=1。 - 只影响录制输出(不改变回放)。
- QoS 便利性:
- 如果使用了
--repeat-transient-local且该话题没有显式 QoS 覆盖,recorder 应为该订阅请求 Transient Local 持久性(以便实际接收到缓存的样本),与 ROS 文档一致。
- 如果使用了
架构:双缓存方法
Section titled “架构:双缓存方法”提议将 transient-local 消息存储在单独的 TransientLocalMessagesCache 中,与现有的 MessageCache 或 CircularMessageCache 并存。这种分离确保 transient-local 消息独立于常规消息流量被保留。
关键原则:
-
独立存储:
TransientLocalMessagesCache维护按话题的 FIFO 队列,深度由--repeat-transient-local选项的用户指定值决定。 -
双写路径: 来自 transient-local 话题的消息同时写入:
- 主缓存(
MessageCache或CircularMessageCache)——保留原始顺序。 TransientLocalMessagesCache——保证不受原始缓存大小约束和旧消息驱逐影响而保留。
- 主缓存(
-
时间戳调整: 在拆分或快照时前置 transient-local 消息时,覆盖时间戳以避免回放间隙:
- 拆分模式:设置为上一个 bag 文件中最后一条消息的时间戳(T_last)。
- 快照模式:设置为快照缓冲区范围内的最早时间戳(T_earliest)。
-
快照模式下去重: 由于消息可能同时存在于两个缓存中,去重会在写入存储前移除重复消息。
TransientLocalMessagesCache 的设计
Section titled “TransientLocalMessagesCache 的设计”目的: 存储每个 transient-local 话题最近的 N 条消息,不受主缓存压力和大小约束的影响。
所需 API:
add_topic(topic_name, queue_depth)- 以指定的队列深度注册话题。remove_topic(topic_name)- 取消注册话题的跟踪。push(topic_name, message)- 向话题队列添加消息(满时驱逐最旧的)。get_messages_sorted_by_timestamp()- 检索所有按 recv_timestamp 排序的缓存消息。clear()- 移除所有缓存消息。size()- 返回缓存消息的总数。
线程安全:
- 必须支持
Recorder的并发写入以及拆分/快照期间Writer的读取。 - 使用互斥锁或无锁结构确保安全的并发访问,且不会显著降低性能。
- 确保用于写入的消息检索是一致的,且不干扰正在进行的录制。
- 设计为最小争用,因为写入预计不频繁(仅针对 transient-local 话题),读取只在拆分/快照期间发生。
拆分模式工作流程
Section titled “拆分模式工作流程”流程:
- 关闭当前存储文件。
- 打开新的存储文件。
- 从
TransientLocalMessagesCache检索 transient-local 消息。 - 将时间戳调整为 T_last(上一个 bag 文件中最后一条消息的时间戳)。
- 将调整后的消息写入新存储文件。
- 恢复正常录制。
时间戳调整:
T_last= 写入上一个 bag 文件的最后一条消息的时间戳(receive_timestamp和send_timestamp)。- 对于每条 transient-local 消息:
- 将
receive_timestamp和send_timestamp都覆盖为T_last中对应的时间戳。这确保回放时,player 看到这些消息就像在拆分时发布的一样,防止原始时间戳在遥远的过去时回放出现大的间隙。
- 将
消息去重:
- 不需要去重,因为 transient-local 消息在新录制继续之前就被预先写入,所以它们自然只会出现一次。
示例:
Split at T=1000.5s, /tf_static originally at T=0.1sAfter prepend: /tf_static written with T=1000.5sNext message at T=1000.6sPlayback: continuous timeline without hours-long gap带去重的快照模式工作流程
Section titled “带去重的快照模式工作流程”挑战: 相同的 transient-local 消息可能存在于两个缓存中:
CircularMessageCache:如果在窗口内收到,保留原始消息顺序。TransientLocalMessagesCache:保证不受缓存压力和旧消息驱逐影响而保留。这可能导致快照输出中出现重复。
流程:
在 CacheConsumer::exec_consuming() 方法中,当快照模式激活且数据就绪时:
- 交换
CircularMessageCache的缓冲区以获得缓存消费者缓冲区。 - 从
TransientLocalMessagesCache检索 transient-local 消息。 - 从缓存消费者缓冲区确定快照范围 [
T_earliest,T_latest]。 - 将早期 transient-local 消息的时间戳调整为
T_earliest。 - 按时间戳范围对从
TransientLocalMessagesCache返回的消息进行去重。 去重确保确定性的输出:同一条消息不会出现两次。 - 合并两个缓存中已排序的消息(合并步骤可以优化并与去重阶段结合,因为两个集合都已按时间戳排序)。
- 将去重后的序列写入存储。
时间戳调整:
T_earliest= 消费者缓冲区中第一条消息的时间戳(receive_timestamp和send_timestamp),即本次快照期间将写入存储的具有最早recv_timestamp的消息。- 对于
msg::recv_timestamp<T_earliest::recv_timestamp的 transient-local 消息:- 将
msg::recv_timestamp覆盖为T_earliest::recv_timestamp - 将
msg::send_timestamp覆盖为T_earliest::send_timestamp
- 将
消息去重: 去重标准:
recv_timestamp>=T_earliest::recv_timestamp的消息被认为是消费者缓冲区中消息的重复,应从与消费者缓冲区的合并中排除。recv_timestamp<T_earliest::recv_timestamp的消息被认为是唯一的,应包含在合并中,但需要调整时间戳,以确保它们出现在快照输出的开头。
集成点和所需 API 变更
Section titled “集成点和所需 API 变更”rosbag2_storage(RecordOptions 扩展):
- 添加字段:
std::unordered_map<std::string, size_t> repeat_transient_local_messages- 将话题名称映射到用于消息保留的队列深度。
- 空映射禁用该功能(保持当前行为)。
ros2bag CLI:
- 添加
add_repeat_transient_local_arg函数来解析--repeat-transient-local参数。 - 用话题名称和深度填充
RecordOptions.repeat_transient_local_messages映射。 - 通过现有的
rosbag2_transport::Recorder构造函数传递配置(RecordOptions)(构造函数不需要 API 变更)。
rosbag2_transport(Recorder 修改):
- 对于每个新订阅,如果话题在
repeat_transient_local_messages中,调用 writer 的create_transient_local_topic()而不是常规的create_topic()。 - 不与
TransientLocalMessagesCache直接交互——所有缓存都在Writer内部处理。 - 消息回调保持不变:简单地调用
writer_->write(message)。 - Writer 在内部处理到存储和 transient-local 缓存的双重写入。
rosbag2_cpp(新的 TransientLocalMessagesCache 类):
- 位置:
rosbag2_cpp/{include,src}/cache/transient_local_messages_cache.{hpp,cpp} - 实现按话题的 FIFO 队列,深度可配置。
- 线程安全,支持并发生产者(Writer 写入路径)和
CacheConsumer。 - 返回按接收时间戳排序的消息,以获得确定性输出。
- 由
SequentialWriter拥有和管理。
rosbag2_cpp(BaseWriterInterface API 扩展):
-
向 writer 接口添加新的虚方法:
virtual void create_transient_local_topic(const rosbag2_storage::TopicMetadata & topic_with_type,size_t num_last_messages) = 0;virtual void create_transient_local_topic(const rosbag2_storage::TopicMetadata & topic_with_type,size_t num_last_messages,const rosbag2_storage::MessageDefinition & message_definition) = 0; -
这些方法注册用于 transient-local 消息缓存的话题,同时在下层存储中创建该话题。
-
num_last_messages参数指定该话题的队列深度。
rosbag2_cpp(SequentialWriter 修改):
- 添加私有成员:
std::shared_ptr<TransientLocalMessagesCache> transient_local_cache_ - 实现
create_transient_local_topic()方法:- 以指定的队列深度在
transient_local_cache_中注册话题。 - 转发到常规的
create_topic()以在存储中注册。
- 以指定的队列深度在
- 修改
write()方法:- 写入存储后,检查该话题是否注册为 transient-local。如果是,也将消息推入
transient_local_cache_。
- 写入存储后,检查该话题是否注册为 transient-local。如果是,也将消息推入
- 修改
split_bagfile():- 如果不在快照模式。打开新的存储文件后,调用内部
prepend_transient_local_messages()。 prepend_transient_local_messages()方法从缓存检索消息,将时间戳调整为T_last,写入存储。
- 如果不在快照模式。打开新的存储文件后,调用内部
rosbag2_cpp(CacheConsumer 修改):
- 在构造函数或通过 setter 接受
TransientLocalMessagesCache引用。 - 在快照模式:
- 消费前检索并合并 transient-local 消息。可以通过检查底层缓存类型(
CircularMessageCache)或通过 Recorder 传入的显式标志来确定快照模式。 - 合并两个缓存的消息时执行时间戳调整和去重,以确保快照输出中没有重复。
- 消费前检索并合并 transient-local 消息。可以通过检查底层缓存类型(
工作流程总结:
- 用户在命令行指定
--repeat-transient-local /topic=N。 - CLI 解析器填充
RecordOptions.repeat_transient_local_messages["/topic"] = N。 - Recorder 在设置期间调用
writer_->create_transient_local_topic("/topic", N)。 SequentialWriter在存储中创建话题,并在内部缓存中以深度 N 注册它。- 每次对该话题调用
writer_->write(msg)时:- 消息写入存储(正常路径)。
- 消息也推入内部
TransientLocalMessagesCache(自动双重写入)。
- 拆分或快照时:
- SequentialWriter 从内部缓存检索消息。
- 调整时间戳,必要时去重。
- 写入新的存储文件。
- Recorder 对缓存实现的细节无感知。
- 通过
--repeat-transient-local标志的可选功能保留了默认的”录制实际发生的事”行为。 - 基于可观察的 bag 文件时间线进行时间戳调整。
- 按快照缓冲区的时间跨度去重确保确定性输出:同一条消息不会出现两次。
- 写入前按接收时间戳排序确保回放期间不会出现乱序消息。
无效配置:
--repeat-transient-local中的话题未被录制:记录警告,继续录制其他话题。- 队列深度
<= 0:以错误消息拒绝。
内存压力:
- 受每个话题用户指定深度的限制,不可能无界增长。
QoS 不匹配:
- 指定了
--repeat-transient-local但发布者不提供 Transient Local:记录警告。
CPU 开销: 双重写入每条消息增加一次共享指针拷贝;快照去重和合并增加 O(M log M) 的排序,其中 M 是缓冲区大小(对于典型的 < 100k 条消息的大小可以忽略不计)。由于两个缓存都已预排序,可以在合并步骤中降低到 O(M)。
磁盘 I/O: 拆分前置只增加最少的写入(每次拆分每条 transient-local 消息一次)。
实施计划(rolling)
Section titled “实施计划(rolling)”阶段 1:核心基础设施
Section titled “阶段 1:核心基础设施”- 在
rosbag2_cpp/cache中实现TransientLocalMessagesCache类- 按话题的 FIFO 队列,深度可配置。
- 线程安全的 push 和检索操作。
- 消息按时间戳排序以获得确定性输出。
- 扩展
rosbag2_storage中的RecordOptions结构体。- 添加
repeat_transient_local_messages映射字段。
- 添加
- 单元测试:
- 验证按话题的队列行为和容量时的驱逐。
- 验证并发访问下的线程安全。
- 验证检索消息的基于时间戳的排序。
阶段 2:Writer API 扩展
Section titled “阶段 2:Writer API 扩展”- 用新的虚方法扩展
BaseWriterInterface:- 带和不带消息定义的
create_transient_local_topic()。
- 带和不带消息定义的
- 修改
SequentialWriter实现:- 添加私有
transient_local_cache_成员。 - 实现
create_transient_local_topic()以注册话题并转发到存储。 - 修改
write(),当话题是 transient-local 时向缓存双重写入消息。
- 添加私有
- 单元测试:
- 验证
create_transient_local_topic()正确注册话题。 - 验证
write()方法中的双重写入行为。 - 验证缓存中每个话题的消息数量正确。
- 验证
阶段 3:拆分模式集成
Section titled “阶段 3:拆分模式集成”- 在
SequentialWriter中实现prepend_transient_local_messages()。- 打开新存储后从缓存检索消息。
- 将时间戳调整为 T_last(上一个 bag 的最后一条消息)。
- 将调整后的消息写入新的存储文件。
- 修改
split_bagfile(),在非快照模式时调用前置逻辑。 - 集成测试:
- 带 transient-local 话题的端到端拆分录制。
- 验证每个拆分文件开头都前置了消息。
- 验证时间戳调整为
T_last。 - 验证输出中没有重复。
阶段 4:快照模式集成
Section titled “阶段 4:快照模式集成”- 修改
CacheConsumer::exec_consuming()以支持快照模式:- 交换 CircularMessageCache 缓冲区后,检索 transient-local 消息。
- 从消费者缓冲区确定快照范围 [
T_earliest,T_latest]。 - 为早于
T_earliest的消息调整时间戳。 - 按时间戳范围去重消息。
- 合并两个缓存中预排序的序列。
- 将合并后的序列传递给消费回调。
- 添加检测快照模式的逻辑(检查缓存类型或显式标志)。
- 集成测试:
- 带 transient-local 话题的端到端快照录制。
- 验证即使在缓存压力下消息仍然存在。
- 验证早期消息的时间戳调整为 T_earliest。
- 验证消息同时出现在两个缓存中时的去重。
阶段 5:Recorder 集成
Section titled “阶段 5:Recorder 集成”- 扩展
Recorder以从选项中解析repeat_transient_local_messages。 - 在订阅设置期间对每个 transient-local 话题:
- 调用
writer_->create_transient_local_topic()而不是常规的create_topic()。
- 调用
- 消息回调不需要更改(双重写入由 Writer 处理)。
- 集成测试:
- 验证 recorder 初始化期间正确的 API 调用。
- 验证 transient-local 话题以正确的深度注册。
- 多个 transient-local 话题的端到端录制。
阶段 6:CLI 和文档
Section titled “阶段 6:CLI 和文档”- 添加
--repeat-transient-local解析。 - 使用示例和最佳实践更新用户文档 README.md 文件。
单元测试(rosbag2_cpp):
-
TransientLocalMessagesCache:- 测试带深度限制的按话题队列。
- 测试满时的最旧消息驱逐。
- 测试线程安全的并发 push 和检索。
- 测试输出的基于时间戳的排序。
- 测试空缓存行为。
- 测试无效配置(深度
<= 0)。
-
SequentialWriter:- 测试
create_transient_local_topic()注册。 - 测试到存储和缓存的双重写入。
- 测试拆分时的前置逻辑。
- 测试时间戳调整为
T_last。 - 测试快照模式下的合并和去重。
- 测试时间戳调整为
T_earliest。
- 测试
集成测试(rosbag2_transport):
-
拆分模式场景:
- 使用较小的
max_bagfile_size录制以触发拆分。 - 在 T=0.1s 发布
/tf_static,在 T=1000.5s 触发拆分。 - 验证每个拆分文件都以
T_last时间戳的/tf_static开头。 - 验证输出中没有重复。
- 在第一次拆分后重新发布
/tf_static,验证它不会出现在下一个拆分文件中。 - 测试多个不同深度的 transient-local 话题。
- 使用较小的
-
快照模式场景:
- 使用较小的
max_cache_size录制以触发缓存压力。 - 在 T=0.1s 发布
/tf_static,用常规消息使缓存溢出。 - 触发快照,验证尽管被驱逐,
/tf_static仍然存在。 - 验证早期消息的时间戳调整为
T_earliest。 - 测试去重:
- 发布
/tf_static,在快照窗口内重新发布。 - 验证输出中没有重复。
- 验证重复消息保留
CircularMessageCache时间戳。
- 发布
- 测试只在快照窗口之前发布的
/tf_static:- 验证时间戳调整为
T_earliest。 - 验证消息恰好出现一次。
- 验证时间戳调整为
- 使用较小的
-
QoS 协商:
- 测试 recorder 自动为 transient-local 话题请求 Transient Local。
- 测试显式 QoS 覆盖优先。
- 测试发布者不提供 Transient Local 时记录警告。
-
边界情况:
- 空缓存(拆分/快照前未收到消息)。
--repeat-transient-local中的话题从未被发布。- 队列深度为 1 与更大深度。
- 多个深度混合的 transient-local 话题。
- transient-local 消息大于缓存容量。
考虑的备选方案
Section titled “考虑的备选方案”带优先级队列的单一缓存: 被拒绝,会使驱逐逻辑复杂化并可能饿死常规消息。
CircularMessageCache 中的百分比预留: 无法保证保留或提供按话题控制。
只特殊处理 /tf_static: 对自定义 transient-local 话题缺乏通用性。
不调整时间戳: 会产生不可接受的回放间隙(player 等待数小时)。
只把 transient-local 消息存储在 circular_message_cache 中: 缺乏对常规 bag 拆分模式的支持。
只存储在 TransientLocalMessagesCache 中: 破坏快照模式保留原始消息顺序的能力。
-
CLI 选项名应该用
--repeat-tl还是--repeat-transient-local还是--repeat-latched?--repeat-latched对 ROS 1 用户很熟悉,但对 ROS 2 用户可能造成困惑,因为 “latched” 不是 ROS 2 QoS 中使用的术语。而且可能不清楚它适用于录制而不是回放。--repeat-tl更简洁,与 ROS 1 的--repeat-latched一致,但描述性较差。--repeat-transient-local更明确但冗长。- 提议:
--repeat-transient-local作为更简洁和明确的选项。用户对--repeat-tl感到困惑,不明白它是什么意思。
-
是否应该在
--repeat-transient-local中支持正则表达式模式(例如/tf_*)?- 根据用户反馈推迟到未来工作。初始实现专注于显式话题名称,以保持简单和清晰。
-
是否应该提供禁用时间戳调整的选项?
- 不。在所有实际场景中,它对于可用的回放都是必要的。
-
是否应该记录去重统计?
- 是的,在测试期间添加 DEBUG 级别日志以提高透明度。
-
是否应该支持不同的时间戳调整策略?
- 推迟。当前策略覆盖所有已知用例。