Skip to content

设计:支持在 Rosbag2 recorder 中重复 transient-local("锁存")消息

作者:Michael Orlov

一些 ROS 2 话题使用 Transient Local 持久性(ROS 1 的”锁存/latched”),使迟到的订阅者也能接收到已发布的数据,但 Rosbag2 用户仍面临两个实际缺口:

  1. 拆分 bag:消费者通常期望每个拆分文件的开头有”锁存”状态(例如 /tf_static、/map),就像 ROS 1 的 --repeat-latched。Rosbag2 目前不提供这个功能。
  2. 快照模式:当循环缓存覆盖早期的 transient-local 消息时,快照可能丢失它们(例如缓存压力下 /tf_static 丢失)。

ROS 2 的”锁存可见性”要求发布者和订阅者都使用 Transient Local 持久性。 当默认值不匹配时,Rosbag2 支持通过 QoS 覆盖文件来处理录制/回放。

  • 为 recorder 拆分和快照提供可选的”重复锁存”体验。
  • 修复快照模式在缓存压力下对选定的 transient-local 话题的驱逐问题。
  • 保持默认行为的确定性:“重复”必须被显式请求。也就是说,transient-local 消息默认不会被重复,否则会引入录制时并未真正发布的额外消息。
  • 自动检测所有 transient-local 话题并默认重复它们。

添加 ros2 bag record --repeat-transient-local 并支持按话题指定消息数量,例如:

Terminal window
$ 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 文档一致。

提议将 transient-local 消息存储在单独的 TransientLocalMessagesCache 中,与现有的 MessageCache 或 CircularMessageCache 并存。这种分离确保 transient-local 消息独立于常规消息流量被保留。

关键原则:

  1. 独立存储: TransientLocalMessagesCache 维护按话题的 FIFO 队列,深度由 --repeat-transient-local 选项的用户指定值决定。

  2. 双写路径: 来自 transient-local 话题的消息同时写入:

    • 主缓存(MessageCache 或 CircularMessageCache)——保留原始顺序。
    • TransientLocalMessagesCache——保证不受原始缓存大小约束和旧消息驱逐影响而保留。
  3. 时间戳调整: 在拆分或快照时前置 transient-local 消息时,覆盖时间戳以避免回放间隙:

    • 拆分模式:设置为上一个 bag 文件中最后一条消息的时间戳(T_last)。
    • 快照模式:设置为快照缓冲区范围内的最早时间戳(T_earliest)。
  4. 快照模式下去重: 由于消息可能同时存在于两个缓存中,去重会在写入存储前移除重复消息。

目的: 存储每个 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 话题),读取只在拆分/快照期间发生。

流程:

  1. 关闭当前存储文件。
  2. 打开新的存储文件。
  3. 从 TransientLocalMessagesCache 检索 transient-local 消息。
  4. 将时间戳调整为 T_last(上一个 bag 文件中最后一条消息的时间戳)。
  5. 将调整后的消息写入新存储文件。
  6. 恢复正常录制。

时间戳调整:

  • 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.1s
After prepend: /tf_static written with T=1000.5s
Next message at T=1000.6s
Playback: continuous timeline without hours-long gap

挑战: 相同的 transient-local 消息可能存在于两个缓存中:

  • CircularMessageCache:如果在窗口内收到,保留原始消息顺序。
  • TransientLocalMessagesCache:保证不受缓存压力和旧消息驱逐影响而保留。这可能导致快照输出中出现重复。

流程:

在 CacheConsumer::exec_consuming() 方法中,当快照模式激活且数据就绪时:

  1. 交换 CircularMessageCache 的缓冲区以获得缓存消费者缓冲区。
  2. 从 TransientLocalMessagesCache 检索 transient-local 消息。
  3. 从缓存消费者缓冲区确定快照范围 [T_earliest, T_latest]。
  4. 将早期 transient-local 消息的时间戳调整为 T_earliest。
  5. 按时间戳范围对从 TransientLocalMessagesCache 返回的消息进行去重。 去重确保确定性的输出:同一条消息不会出现两次。
  6. 合并两个缓存中已排序的消息(合并步骤可以优化并与去重阶段结合,因为两个集合都已按时间戳排序)。
  7. 将去重后的序列写入存储。

时间戳调整:

  • 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 的消息被认为是唯一的,应包含在合并中,但需要调整时间戳,以确保它们出现在快照输出的开头。

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_。
  • 修改 split_bagfile():
    • 如果不在快照模式。打开新的存储文件后,调用内部 prepend_transient_local_messages()。
    • prepend_transient_local_messages() 方法从缓存检索消息,将时间戳调整为 T_last,写入存储。

rosbag2_cpp(CacheConsumer 修改):

  • 在构造函数或通过 setter 接受 TransientLocalMessagesCache 引用。
  • 在快照模式:
    • 消费前检索并合并 transient-local 消息。可以通过检查底层缓存类型(CircularMessageCache)或通过 Recorder 传入的显式标志来确定快照模式。
    • 合并两个缓存的消息时执行时间戳调整和去重,以确保快照输出中没有重复。

工作流程总结:

  1. 用户在命令行指定 --repeat-transient-local /topic=N。
  2. CLI 解析器填充 RecordOptions.repeat_transient_local_messages["/topic"] = N。
  3. Recorder 在设置期间调用 writer_->create_transient_local_topic("/topic", N)。
  4. SequentialWriter 在存储中创建话题,并在内部缓存中以深度 N 注册它。
  5. 每次对该话题调用 writer_->write(msg) 时:
    • 消息写入存储(正常路径)。
    • 消息也推入内部 TransientLocalMessagesCache(自动双重写入)。
  6. 拆分或快照时:
    • SequentialWriter 从内部缓存检索消息。
    • 调整时间戳,必要时去重。
    • 写入新的存储文件。
  7. 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 消息一次)。

  • 在 rosbag2_cpp/cache 中实现 TransientLocalMessagesCache 类
    • 按话题的 FIFO 队列,深度可配置。
    • 线程安全的 push 和检索操作。
    • 消息按时间戳排序以获得确定性输出。
  • 扩展 rosbag2_storage 中的 RecordOptions 结构体。
    • 添加 repeat_transient_local_messages 映射字段。
  • 单元测试:
    • 验证按话题的队列行为和容量时的驱逐。
    • 验证并发访问下的线程安全。
    • 验证检索消息的基于时间戳的排序。
  • 用新的虚方法扩展 BaseWriterInterface:
    • 带和不带消息定义的 create_transient_local_topic()。
  • 修改 SequentialWriter 实现:
    • 添加私有 transient_local_cache_ 成员。
    • 实现 create_transient_local_topic() 以注册话题并转发到存储。
    • 修改 write(),当话题是 transient-local 时向缓存双重写入消息。
  • 单元测试:
    • 验证 create_transient_local_topic() 正确注册话题。
    • 验证 write() 方法中的双重写入行为。
    • 验证缓存中每个话题的消息数量正确。
  • 在 SequentialWriter 中实现 prepend_transient_local_messages()。
    • 打开新存储后从缓存检索消息。
    • 将时间戳调整为 T_last(上一个 bag 的最后一条消息)。
    • 将调整后的消息写入新的存储文件。
  • 修改 split_bagfile(),在非快照模式时调用前置逻辑。
  • 集成测试:
    • 带 transient-local 话题的端到端拆分录制。
    • 验证每个拆分文件开头都前置了消息。
    • 验证时间戳调整为 T_last。
    • 验证输出中没有重复。
  • 修改 CacheConsumer::exec_consuming() 以支持快照模式:
    • 交换 CircularMessageCache 缓冲区后,检索 transient-local 消息。
    • 从消费者缓冲区确定快照范围 [T_earliest, T_latest]。
    • 为早于 T_earliest 的消息调整时间戳。
    • 按时间戳范围去重消息。
    • 合并两个缓存中预排序的序列。
    • 将合并后的序列传递给消费回调。
  • 添加检测快照模式的逻辑(检查缓存类型或显式标志)。
  • 集成测试:
    • 带 transient-local 话题的端到端快照录制。
    • 验证即使在缓存压力下消息仍然存在。
    • 验证早期消息的时间戳调整为 T_earliest。
    • 验证消息同时出现在两个缓存中时的去重。
  • 扩展 Recorder 以从选项中解析 repeat_transient_local_messages。
  • 在订阅设置期间对每个 transient-local 话题:
    • 调用 writer_->create_transient_local_topic() 而不是常规的 create_topic()。
  • 消息回调不需要更改(双重写入由 Writer 处理)。
  • 集成测试:
    • 验证 recorder 初始化期间正确的 API 调用。
    • 验证 transient-local 话题以正确的深度注册。
    • 多个 transient-local 话题的端到端录制。
  • 添加 --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 消息大于缓存容量。

带优先级队列的单一缓存: 被拒绝,会使驱逐逻辑复杂化并可能饿死常规消息。

CircularMessageCache 中的百分比预留: 无法保证保留或提供按话题控制。

只特殊处理 /tf_static: 对自定义 transient-local 话题缺乏通用性。

不调整时间戳: 会产生不可接受的回放间隙(player 等待数小时)。

只把 transient-local 消息存储在 circular_message_cache 中: 缺乏对常规 bag 拆分模式的支持。

只存储在 TransientLocalMessagesCache 中: 破坏快照模式保留原始消息顺序的能力。

  1. 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 感到困惑,不明白它是什么意思。
  2. 是否应该在 --repeat-transient-local 中支持正则表达式模式(例如 /tf_*)?

    • 根据用户反馈推迟到未来工作。初始实现专注于显式话题名称,以保持简单和清晰。
  3. 是否应该提供禁用时间戳调整的选项?

    • 不。在所有实际场景中,它对于可用的回放都是必要的。
  4. 是否应该记录去重统计?

    • 是的,在测试期间添加 DEBUG 级别日志以提高透明度。
  5. 是否应该支持不同的时间戳调整策略?

    • 推迟。当前策略覆盖所有已知用例。