rosbag2 源码详细分析

rosbag2 源码详细分析

工作区路径:/home/cp/work2/ros2Learn/ros2_humble/src/ros2/rosbag2
版本:0.15.16(Humble),子包 19 个,许可证 Apache 2.0

rosbag2 是 ROS 2 录制与回放 的完整实现栈:可插拔存储后端(SQLite3 / MCAP)、消息/文件级压缩、序列化格式转换、C++/Python API,以及通过 ros2 bag 暴露的 CLI。设计文档见 ros2/design rosbags

范围说明:本仓库不含 ROS 1 bag 读取插件(rosbag2_bag_v2 在独立仓库);Humble 默认安装 sqlite3 + zstd 插件,MCAP 需额外安装 rosbag2_storage_mcap


1. 总体认识

1.1 核心职责

能力 实现位置 说明
存储抽象 rosbag2_storage ReadOnlyInterface / ReadWriteInterface,pluginlib 加载
默认 SQLite3 rosbag2_storage_default_plugins .db3 文件,默认 storage_id=sqlite3
MCAP 格式 rosbag2_storage_mcap Foxglove 生态,单文件/索引友好
读写 API rosbag2_cpp Reader / WriterSequentialReader/Writer
与 ROS 图集成 rosbag2_transport Recorder / Player,GenericSubscription/Publisher
压缩 rosbag2_compression + _zstd MESSAGE / FILE 两种模式
Python 绑定 rosbag2_py pybind11 暴露 Recorder/Player/convert 等
CLI ros2bag 注册到 ros2cli.commandbag 命令
消息/服务定义 rosbag2_interfaces 回放控制服务、split 事件

1.2 在 ROS 2 栈中的位置

用户ros2bagrosbag2_pyrosbag2_transportrosbag2_cpprosbag2_storage + pluginsrosbag2_compressionROS 2 运行时ros2 bag record / playVerbExtension\nrecord/play/info/convert/reindex/listpybind11\nRecorder / Player / bag_rewriteRecorder\nGenericSubscriptionPlayer\nGenericPublisher + TimeControllerClockReaderWriterFactoryWriter → SequentialWriterReader → SequentialReaderMessageCache\n双缓冲StorageFactory\npluginlibsqlite3mcapSequentialCompressionWriterSequentialCompressionReaderzstd pluginrclcpp / rmwDDS 图
对比项 录制 (Recorder) 回放 (Player)
ROS 接口 GenericSubscription 订阅 live topic GenericPublisher 重放历史消息
时间源 系统时钟或 /clock(sim time) TimeControllerClock 控制回放节奏
写入/读取 Writer::write(SerializedMessage) Reader::read_next()
默认存储 SequentialWriter + sqlite3 SequentialReader,压缩 bag 自动选 SequentialCompressionReader
发现 周期性 topic discovery(可关闭) 无 discovery,topic 来自 bag metadata

1.3 命令行分层(与 ros2cli 集成)

1
2
3
4
5
6
7
8
ros2                          ← ros2cli 主入口
└── bag ← entry point: ros2bag.command.bag:BagCommand
├── record ← ros2bag.verb.record:RecordVerb
├── play ← ros2bag.verb.play:PlayVerb
├── info ← ros2bag.verb.info:InfoVerb
├── convert ← ros2bag.verb.convert:ConvertVerb
├── reindex ← ros2bag.verb.reindex:ReindexVerb
└── list ← ros2bag.verb.list:ListVerb

ros2bag/setup.py 通过 setuptools entry_points 注册;BagCommand 使用 add_subparsers_on_demand 按需加载 verb 插件,与 ros2 topic 等命令模式一致。


2. 仓库结构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
rosbag2/
├── README.md # 用户文档(record/play/split/compression/convert)
├── rosbag2/ # ★ metapackage 0.15.16
├── rosbag2_storage/ # ★ 存储抽象 (~2500 行 C++)
│ ├── include/rosbag2_storage/
│ │ ├── storage_interfaces/ # ReadOnly / ReadWrite / BaseWrite
│ │ ├── serialized_bag_message.hpp
│ │ ├── bag_metadata.hpp
│ │ ├── topic_metadata.hpp
│ │ ├── storage_options.hpp
│ │ └── storage_factory.hpp
│ └── src/.../storage_factory_impl.hpp # pluginlib 加载逻辑
├── rosbag2_storage_default_plugins/ # SQLite3 插件
├── rosbag2_storage_mcap/ # MCAP 插件
├── rosbag2_cpp/ # ★ 核心 API (~9400 行)
│ ├── writer.hpp / reader.hpp
│ ├── writers/sequential_writer.cpp
│ ├── readers/sequential_reader.cpp
│ ├── cache/message_cache.hpp # 双缓冲写缓存
│ └── reindexer.cpp
├── rosbag2_compression/ # 压缩读写包装 (~3100 行)
├── rosbag2_compression_zstd/ # zstd compressor 插件
├── rosbag2_transport/ # ★ Recorder/Player (~8600 行)
│ ├── recorder.cpp / player.cpp
│ ├── bag_rewrite.cpp
│ └── reader_writer_factory.cpp
├── rosbag2_py/ # Python 绑定
├── ros2bag/ # CLI (Python)
├── rosbag2_interfaces/ # msg/srv 定义
├── sqlite3_vendor / zstd_vendor / mcap_vendor / shared_queues_vendor
├── rosbag2_tests / rosbag2_test_common
└── rosbag2_performance/ # 性能基准

2.1 子包一览

包名 类型 职责
rosbag2 metapackage 聚合依赖,默认带 sqlite3 + zstd
rosbag2_storage 存储接口、metadata IO、StorageFactory
rosbag2_storage_default_plugins 插件 sqlite3SqliteStorage
rosbag2_storage_mcap 插件 mcap 存储
rosbag2_cpp Reader/Writer、顺序读写、Reindexer、Info
rosbag2_transport Recorder、Player、bag_rewrite、工厂
rosbag2_compression SequentialCompressionReader/Writer
rosbag2_compression_zstd 插件 zstd 压缩器
rosbag2_py pybind11 绑定
ros2bag CLI ros2 bag 命令
rosbag2_interfaces 接口 WriteSplitEvent、Playback 控制 srv
sqlite3_vendor vendor 打包 sqlite3
zstd_vendor vendor 打包 zstd
mcap_vendor vendor 打包 mcap C++ 库
shared_queues_vendor vendor moodycamel 无锁队列(Player 预读)
rosbag2_tests 测试 集成测试
rosbag2_test_common 测试 公共测试工具
rosbag2_storage_mcap_testdata 测试 MCAP 测试数据
rosbag2_performance_benchmarking 工具 性能对比

3. 三层插件体系

rosbag2 的可扩展性来自 三类独立插件,均通过 pluginlib 或工厂模式注册:

3.1 存储插件 (Storage)

基类rosbag2_storage::storage_interfaces::ReadWriteInterface(读写合一)
只读变体ReadOnlyInterface(部分插件可只读打开)

加载路径storage_factory_impl.hpp):

  1. storage_options.storage_id 非空 → 按 ID 加载对应 class(如 sqlite3mcap
  2. 若为空且为 → 遍历已注册插件,逐个 open() 尝试,首个成功者选用
  3. 读模式下若 ReadOnly 插件失败 → 回退尝试 ReadWrite 插件以 READ_ONLY 打开
1
2
3
4
<!-- rosbag2_storage_default_plugins/plugin_description.xml -->
<class name="sqlite3"
type="rosbag2_storage_plugins::SqliteStorage"
base_class_type="rosbag2_storage::storage_interfaces::ReadWriteInterface"/>

SQLite3 实现要点sqlite_storage.cpp):

  • 表结构存储 topic metadata 与 serialized message blob
  • 支持 storage_config_uri YAML 配置 SQLite pragma(read/write 分节)
  • get_minimum_split_file_size() 返回最小 split 阈值(默认约 86KB,与 CLI -b 说明一致)
  • 支持 IOFlag::READ_ONLY / READ_WRITE / APPEND

MCAProsbag2_storage_mcap):面向 robotics 日志的单文件格式,索引与 Foxglove Studio 兼容;通过 --storage mcap 选用。

3.2 序列化格式转换 (Serialization Converter)

目的:bag 内存储的序列化格式可与录制时 RMW 格式不同(如 cdr ↔ protobuf 实验性路径)。

  • ConverterOptionsinput_serialization_format / output_serialization_format
  • SerializationFormatConverterFactoryInterface 加载 *_converter 插件
  • SequentialWriter::open() 若输入/输出格式不同则创建 converter_,在 write() 路径转换

Recorder 打开 writer 时传入:

1
2
writer_->open(storage_options_,
{rmw_get_serialization_format(), record_options_.rmw_serialization_format});

即:输入为当前 RMW 格式,输出为 CLI --serialization-format 指定格式(默认同 RMW)。

3.3 压缩插件 (Compression)

配置compression_options.hpp):

字段 含义
compression_format zstd(来自 rosbag2_compression_zstd
compression_mode NONE / FILE / MESSAGE
compression_queue_size 异步压缩队列深度
compression_threads 压缩线程数(0 → 硬件并发数)

工厂选择reader_writer_factory.cpp):

  • 录制:compression_format 非空 → SequentialCompressionWriter 包装 SequentialWriter

  • 回放:读 metadata.yamlcompression_format 非空 → SequentialCompressionReader

  • MESSAGE 模式:每条消息压缩后写入 storage

  • FILE 模式:关闭 bag 文件时对整文件压缩(split 时按文件压缩)


4. 数据模型

4.1 SerializedBagMessage

单条 bag 消息的最小单元(serialized_bag_message.hpp):

1
2
3
4
5
struct SerializedBagMessage {
std::shared_ptr<rcutils_uint8_array_t> serialized_data;
rcutils_time_point_value_t time_stamp; // 纳秒时间戳
std::string topic_name;
};
  • 不做 ROS 类型反序列化;payload 为 RMW 序列化字节流
  • Recorder 从 rclcpp::SerializedMessage 拷贝 buffer 与时间戳写入

4.2 TopicMetadata

1
2
3
4
5
6
struct TopicMetadata {
std::string name;
std::string type; // 如 sensor_msgs/msg/Image
std::string serialization_format; // 如 cdr
std::string offered_qos_profiles; // YAML 序列化的 QoS 列表
};

回放时 Player 用 metadata 中的 type 创建 GenericPublisher,并按 recorded QoS(或 override)发布。

4.3 BagMetadata 与 metadata.yaml

BagMetadatabag_metadata.hppversion = 5):

字段 说明
storage_identifier sqlite3mcap
relative_file_paths bag 目录内相对路径列表
files[] 每个分片的 path、starting_time、duration、message_count
starting_time / duration / message_count 整 bag 统计
topics_with_message_count[] topic 元数据 + 消息计数
compression_format / compression_mode 压缩信息

MetadataIo 负责读写 bag 目录下的 metadata.yamlSequentialReader 打开时先读 metadata,再按 relative_file_paths 解析分片路径(version < 4 的路径规则不同,见 resolve_relative_paths)。

典型目录结构:

1
2
3
4
my_bag/
├── metadata.yaml
├── my_bag_0.db3 # sqlite3 分片
└── my_bag_0.db3.zstd # FILE 压缩时可能出现

4.4 StorageOptions

录制/回放/convert 共用(storage_options.hpp):

选项 默认 说明
uri bag 目录路径
storage_id 写时默认 sqlite3 存储插件 ID
max_bagfile_size 0 字节,超过则 split
max_bagfile_duration 0 秒,超过则 split
max_cache_size 0 消息缓存条数,0=直写磁盘
storage_preset_profile “” 预设配置名
storage_config_uri “” 插件专用 YAML(如 SQLite pragma)
snapshot_mode false 环形缓冲 + ~/snapshot 服务

5. rosbag2_cpp:读写核心

5.1 Writer 门面

rosbag2_cpp::Writer 持有 BaseWriterInterface 实现(默认 SequentialWriter),提供:

  • open(storage_options, converter_options)
  • create_topic / remove_topic
  • write(SerializedBagMessage)write(SerializedMessage, topic, type, time)
  • take_snapshot()(snapshot 模式)
  • add_event_callbacks(如 split 事件)

线程安全Writer 层有 mutex,防止并发 open/close/write 竞态。

5.2 SequentialWriter 流程

  1. open()storage_factory_->open_read_write() 创建插件实例
  2. init_metadata():初始化 BagMetadata,记录首个相对路径
  3. create_topic():写 topic 行到 storage + 更新 metadata
  4. write()
    • 可选 converter 转换序列化格式
    • max_cache_size > 0MessageCache 双缓冲;否则直接 storage_->write()
    • 更新 duration、message_count
  5. close():刷 cache、写最终 metadata.yaml

Split 逻辑:当单文件 size 或 duration 超阈值,关闭当前 storage、新建分片文件,更新 relative_file_paths,触发 write_split_callback(Recorder 发布 WriteSplitEvent)。

默认常量:

1
static constexpr char const * kDefaultStorageID = "sqlite3";

5.3 MessageCache(双缓冲)

message_cache.hpp 实现 greedy consumer 双缓冲:

  • Producer(subscription 回调)push 到主 buffer
  • Consumer(磁盘写入线程)swap_buffers 取走数据
  • buffer 满时 丢弃 消息并计数(性能告警信号)
  • flush 状态用于 shutdown 时排空

设计目标:避免慢速磁盘阻塞 DDS 回调线程。

5.4 SequentialReader 流程

  1. open():读 metadata → 打开首个 storage 分片
  2. has_next() / read_next():顺序返回 SerializedBagMessage
  3. 当前分片读完 → 自动打开 relative_file_paths 中下一个
  4. 支持 seek(timestamp)set_filter(StorageFilter)(topic 白名单)
  5. 若存在 converter,读出的消息可转换为目标序列化格式

5.5 Reindexer

metadata.yaml 损坏或与 .db3 不一致时,Reindexer::reindex() 扫描目录内数据库文件,重建 metadata(reindexer.cpp)。CLI:ros2 bag reindex


6. rosbag2_transport:录制路径

6.1 Recorder 架构

Recorder 继承 rclcpp::Node,核心成员:

  • Writer(通过 ReaderWriterFactory::make_writer 创建,可能为压缩 writer)
  • subscriptions_GenericSubscription 映射
  • topics_discovery 异步线程(除非 --no-discovery
  • paused_ 原子标志
  • snapshot 模式下的 ~/snapshot 服务

6.2 record() 主流程

1
2
3
4
5
6
7
record()
├─ writer_->open(storage_options, converter_options)
├─ [snapshot_mode] 创建 Snapshot 服务
├─ 注册 write_split 回调 → 发布 WriteSplitEvent
├─ subscribe_topics(初始 topic 集)
├─ [可选] std::async topics_discovery 循环
└─ 等待 spin(CLI 侧 rclpy/py 驱动 executor)

6.3 订阅与写入

关键顺序(避免竞态):

1
2
3
4
5
6
7
void Recorder::subscribe_topic(const TopicMetadata & topic) {
writer_->create_topic(topic); // 必须先注册 topic
auto sub = create_generic_subscription(..., [this](SerializedMessage msg) {
if (!paused_.load())
writer_->write(msg, topic_name, topic_type, get_clock()->now());
});
}

QoS 适配

  • 默认 Rosbag2QoS::adapt_request_to_offers() 根据 publisher 的 offered QoS 选择 subscription QoS
  • --qos-profile-overrides-path 可 per-topic 覆盖
  • discovery 过程中若新 publisher QoS 与已订阅不兼容,会 warn_if_new_qos_for_subscribed_topic

Topic 选择

  • 显式 topic 列表、-a(全部)、-e / -x 正则过滤
  • get_requested_or_available_topics() 合并请求与图发现结果
  • --include-hidden-topics / --include-unpublished-topics 扩展发现范围

Sim timeuse_sim_time 时 Recorder 使用 /clock;在收到首个 clock 前不应写消息(README 警告 time=0 问题)。

键盘控制:空格切换 pause/resume(keyboard_handler)。


7. rosbag2_transport:回放路径

7.1 Player 架构

  • Reader + TimeControllerClock(速率、pause、seek)
  • moodycamel::ReaderWriterQueue 预读队列(read_ahead_queue_size,默认 1000)
  • GenericPublisher per topic
  • 可选 /clock 定时发布(clock_publish_frequency
  • 控制服务:pauseresumetoggle_pausedset_rateseekburstplay_next 等(rosbag2_interfaces

7.2 play() 主流程

1
2
3
4
5
6
7
8
9
10
play()
loop (若 play_options_.loop):
├─ [delay] sleep
├─ reader_->seek(starting_time_); clock_->jump(...)
├─ async load_storage_content() // 后台读 bag 填队列
├─ wait_for_filled_queue()
├─ play_messages_from_queue()
│ └─ clock_->sleep_until(msg.timestamp) // 节奏控制
│ └─ generic_publisher->publish(serialized)
└─ [可选] wait_for_all_acked

时间控制TimeControllerClock 用 steady clock 映射 bag 时间戳,支持 rate 倍速、pauseseek

QoS:默认从 bag metadata 恢复;PlayOptions::topic_qos_profile_overrides 可覆盖。

Remappingtopic_remapping_options 支持回放时改名。


8. bag_rewrite 与 convert

rosbag2_transport::bag_rewrite()bag_rewrite.hpp)是 通用 bag 变换引擎

输入:多个 StorageOptions(多个源 bag)
输出:多个 (StorageOptions, RecordOptions) 对(目标 bag 配置)

能力组合:

场景 实现方式
合并多 bag 多 input → 单 output
格式转换 storage_id(sqlite3 → mcap)
压缩/解压 RecordOptions.compression_*
序列化转换 rmw_serialization_format
过滤 topic output RecordOptions.topics
Split output max_bagfile_size / duration

CLI:ros2 bag convert -i input_options.yaml -o output_options.yaml
Python:rosbag2_py.bag_rewrite()
C++:rosbag2_transport::bag_rewrite()


9. rosbag2_py 与 ros2bag CLI

9.1 pybind11 模块(_transport.cpp 等)

暴露类型:

  • StorageOptions / RecordOptions / PlayOptions
  • Recorder / Player
  • bag_rewrite / Reindexer / Info
  • get_registered_writers() / get_registered_compressors() 等 introspection

RecordVerb 典型用法:

1
2
3
from rosbag2_py import Recorder, RecordOptions, StorageOptions
recorder = Recorder('my_recorder', '', None)
recorder.record(storage_options, record_options)

CLI 解析参数后构造 C++ 对象,在 Python 侧 rclpy.spin 或直接调用阻塞式 record()

9.2 各 verb 职责

Verb 调用链
record RecordVerbrosbag2_py.RecorderRecorder::record
play PlayVerbrosbag2_py.PlayerPlayer::play
info InfoVerbrosbag2_py info API → 读 metadata
convert ConvertVerbbag_rewrite
reindex ReindexVerbReindexer
list ListVerb → 列出已注册 storage/compression 插件

10. rosbag2_interfaces

类型 名称 用途
msg WriteSplitEvent 录制分片通知(closed/opened file)
msg ReadSplitEvent 回放切换分片(内部)
srv Snapshot snapshot 模式触发落盘
srv Pause / Resume / TogglePaused / IsPaused 回放控制
srv SetRate / GetRate 倍速
srv Seek 跳转时间点
srv Burst / PlayNext 单步/突发播放

Recorder 发布 events/write_split;Player 提供 playback 控制服务(默认 namespace 下)。


11. 高级功能摘要

11.1 Bag splitting

  • CLI:-b 字节阈值、-d 秒阈值
  • 二者同时启用时 先触达者 split
  • SequentialWriter 负责切换文件;metadata 记录所有分片

11.2 Snapshot 模式

  • StorageOptions.snapshot_mode = true
  • Writer 维护环形缓冲,仅保留最近 N 消息
  • 调用 ~/snapshot 服务将缓冲 flush 到磁盘

11.3 QoS override 文件

YAML 指定 topic → QoS profile,用于录制订阅或回放发布,解决 recorded QoS 与当前系统不匹配问题。

11.4 性能相关选项

  • max_cache_size:录制写缓存
  • compression_queue_size / compression_threads:压缩流水线
  • read_ahead_queue_size:回放预读
  • SQLite storage_config_uri pragma tuning(如 WAL、synchronous)

12. 与 ROS 1 rosbag 对比

维度 ROS 1 rosbag rosbag2
存储格式 单一 .bag(自定义) 插件化(sqlite3/mcap/…)
消息表示 录制时反序列化再存 默认存 RMW 序列化字节(更高效、格式可转换)
元数据 bag 内嵌 目录 + metadata.yaml + 数据文件
QoS 无 DDS QoS 录制/回放 offered_qos_profiles
压缩 bz2 zstd,MESSAGE/FILE 模式
CLI rosbag record/play ros2 bag record/play/...
工具链 紧耦合 分层:storage / cpp / transport / py
Sim time /clock 同样支持,实现于 Recorder/Player

ROS 1 bag 直接读取需额外 rosbag2_bag_v2 插件(本仓库不含)。


13. 端到端序列图

13.1 录制

SqliteStorageMessageCacheWriterRecorderGenericSubscriptionDDSSqliteStorageMessageCacheWriterRecorderGenericSubscriptionDDSalt[not paused]close 时写 metadata.yamlSerializedMessagecallbackwrite(msg, topic, time)push (if cache enabled)flush batchwrite(SerializedBagMessage)

13.2 回放

DDSGenericPublisherPlayerReadAheadQueueReaderStorageDDSGenericPublisherPlayerReadAheadQueueReaderStorageloop[load_storage_con-tent]loop[play_messages_fr-om_queue]open(uri)read metadata + messagesenqueue messagespop messageclock.sleep_until(ts)publish(serialized)publish

14. 调试与常见问题

现象 排查方向
录不到消息 topic 是否 discovery 到;是否 start_paused;sim time 下是否收到 /clock
QoS 不匹配丢包 对比 ros2 topic info -v 与 bag metadata QoS;尝试 overrides YAML
回放无订阅者 检查 topic remapping、filter、类型名是否一致
metadata 损坏 ros2 bag reindex
转换失败 检查 storage_id、compression、serialization converter 是否安装
性能差 / 丢消息 增大 max_cache_size;检查 MessageCache dropped 计数;磁盘 IO
插件加载失败 ros2 bag list 看注册插件;pluginlib 库路径 / AMENT_PREFIX_PATH

日志:各层使用 ROS_BAG_LOG / RCLCPP_*;storage 层 ROSBAG2_STORAGE_LOG_*

实用命令

1
2
3
4
ros2 bag info my_bag
ros2 bag list # 已注册 storage/compressor
ros2 bag record -a --storage mcap
ros2 bag play my_bag --rate 2.0 --clock 100

15. 源码阅读顺序

  1. 数据模型serialized_bag_message.hpptopic_metadata.hppbag_metadata.hppstorage_options.hpp
  2. 插件加载storage_factory_impl.hppplugin_description.xml(sqlite3)
  3. 写路径writer.hppsequential_writer.cppmessage_cache.hpp
  4. 读路径reader.hppsequential_reader.cpp
  5. 压缩reader_writer_factory.cppsequential_compression_writer.cpp
  6. 录制record_options.hpprecorder.cppsubscribe_topic / record
  7. 回放play_options.hppplayer.cppplay / play_messages_from_queue
  8. CLIros2bag/setup.pyverb/record.pyrosbag2_py/_transport.cpp
  9. 高级bag_rewrite.cppreindexer.cpp

16. 小结

rosbag2 目录实现 ROS 2 可插拔日志栈:底层 storage 插件 持久化序列化字节;rosbag2_cpp 提供顺序读写、缓存与 metadata;rosbag2_transport 连接 rclcpp 图(GenericSubscription/Publisher)与时间控制;rosbag2_compressionros2bag/rosbag2_py 完成压缩与用户接口。

理解任意 bug 或性能问题,通常从 Recorder::subscribe_topicWriter::write → storage 插件(录制)或 Player::playReader::read_next → GenericPublisher(回放)两条主链入手,并核对 metadata.yaml 与 QoS/serialization 配置 是否一致。

文章互动

阅读 --

留言

0 条留言

正在加载留言…