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 / Writer,SequentialReader/Writer |
| 与 ROS 图集成 | rosbag2_transport |
Recorder / Player,GenericSubscription/Publisher |
| 压缩 | rosbag2_compression + _zstd |
MESSAGE / FILE 两种模式 |
| Python 绑定 | rosbag2_py |
pybind11 暴露 Recorder/Player/convert 等 |
| CLI | ros2bag |
注册到 ros2cli.command 的 bag 命令 |
| 消息/服务定义 | rosbag2_interfaces |
回放控制服务、split 事件 |
1.2 在 ROS 2 栈中的位置
| 对比项 | 录制 (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 | ros2 ← ros2cli 主入口 |
ros2bag/setup.py 通过 setuptools entry_points 注册;BagCommand 使用 add_subparsers_on_demand 按需加载 verb 插件,与 ros2 topic 等命令模式一致。
2. 仓库结构
1 | rosbag2/ |
2.1 子包一览
| 包名 | 类型 | 职责 |
|---|---|---|
rosbag2 |
metapackage | 聚合依赖,默认带 sqlite3 + zstd |
rosbag2_storage |
库 | 存储接口、metadata IO、StorageFactory |
rosbag2_storage_default_plugins |
插件 | sqlite3 → SqliteStorage |
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):
- 若
storage_options.storage_id非空 → 按 ID 加载对应 class(如sqlite3、mcap) - 若为空且为 读 → 遍历已注册插件,逐个
open()尝试,首个成功者选用 - 读模式下若 ReadOnly 插件失败 → 回退尝试 ReadWrite 插件以 READ_ONLY 打开
1 | <!-- rosbag2_storage_default_plugins/plugin_description.xml --> |
SQLite3 实现要点(sqlite_storage.cpp):
- 表结构存储 topic metadata 与 serialized message blob
- 支持
storage_config_uriYAML 配置 SQLite pragma(read/write 分节) get_minimum_split_file_size()返回最小 split 阈值(默认约 86KB,与 CLI-b说明一致)- 支持
IOFlag::READ_ONLY/READ_WRITE/APPEND
MCAP(rosbag2_storage_mcap):面向 robotics 日志的单文件格式,索引与 Foxglove Studio 兼容;通过 --storage mcap 选用。
3.2 序列化格式转换 (Serialization Converter)
目的:bag 内存储的序列化格式可与录制时 RMW 格式不同(如 cdr ↔ protobuf 实验性路径)。
ConverterOptions:input_serialization_format/output_serialization_formatSerializationFormatConverterFactoryInterface加载*_converter插件SequentialWriter::open()若输入/输出格式不同则创建converter_,在write()路径转换
Recorder 打开 writer 时传入:
1 | writer_->open(storage_options_, |
即:输入为当前 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.yaml中compression_format非空 →SequentialCompressionReaderMESSAGE 模式:每条消息压缩后写入 storage
FILE 模式:关闭 bag 文件时对整文件压缩(split 时按文件压缩)
4. 数据模型
4.1 SerializedBagMessage
单条 bag 消息的最小单元(serialized_bag_message.hpp):
1 | struct SerializedBagMessage { |
- 不做 ROS 类型反序列化;payload 为 RMW 序列化字节流
- Recorder 从
rclcpp::SerializedMessage拷贝 buffer 与时间戳写入
4.2 TopicMetadata
1 | struct TopicMetadata { |
回放时 Player 用 metadata 中的 type 创建 GenericPublisher,并按 recorded QoS(或 override)发布。
4.3 BagMetadata 与 metadata.yaml
BagMetadata(bag_metadata.hpp,version = 5):
| 字段 | 说明 |
|---|---|
storage_identifier |
如 sqlite3、mcap |
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.yaml;SequentialReader 打开时先读 metadata,再按 relative_file_paths 解析分片路径(version < 4 的路径规则不同,见 resolve_relative_paths)。
典型目录结构:
1 | my_bag/ |
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_topicwrite(SerializedBagMessage)或write(SerializedMessage, topic, type, time)take_snapshot()(snapshot 模式)add_event_callbacks(如 split 事件)
线程安全:Writer 层有 mutex,防止并发 open/close/write 竞态。
5.2 SequentialWriter 流程
open():storage_factory_->open_read_write()创建插件实例init_metadata():初始化BagMetadata,记录首个相对路径create_topic():写 topic 行到 storage + 更新 metadatawrite():- 可选 converter 转换序列化格式
- 若
max_cache_size > 0→MessageCache双缓冲;否则直接storage_->write() - 更新 duration、message_count
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 流程
open():读 metadata → 打开首个 storage 分片has_next()/read_next():顺序返回SerializedBagMessage- 当前分片读完 → 自动打开
relative_file_paths中下一个 - 支持
seek(timestamp)、set_filter(StorageFilter)(topic 白名单) - 若存在 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 | record() |
6.3 订阅与写入
关键顺序(避免竞态):
1 | void Recorder::subscribe_topic(const TopicMetadata & topic) { |
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 time:use_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)GenericPublisherper topic- 可选
/clock定时发布(clock_publish_frequency) - 控制服务:
pause、resume、toggle_paused、set_rate、seek、burst、play_next等(rosbag2_interfaces)
7.2 play() 主流程
1 | play() |
时间控制:TimeControllerClock 用 steady clock 映射 bag 时间戳,支持 rate 倍速、pause、seek。
QoS:默认从 bag metadata 恢复;PlayOptions::topic_qos_profile_overrides 可覆盖。
Remapping:topic_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/PlayOptionsRecorder/Playerbag_rewrite/Reindexer/Infoget_registered_writers()/get_registered_compressors()等 introspection
RecordVerb 典型用法:
1 | from rosbag2_py import Recorder, RecordOptions, StorageOptions |
CLI 解析参数后构造 C++ 对象,在 Python 侧 rclpy.spin 或直接调用阻塞式 record()。
9.2 各 verb 职责
| Verb | 调用链 |
|---|---|
record |
RecordVerb → rosbag2_py.Recorder → Recorder::record |
play |
PlayVerb → rosbag2_py.Player → Player::play |
info |
InfoVerb → rosbag2_py info API → 读 metadata |
convert |
ConvertVerb → bag_rewrite |
reindex |
ReindexVerb → Reindexer |
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_uripragma 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 录制
13.2 回放
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 | ros2 bag info my_bag |
15. 源码阅读顺序
- 数据模型:
serialized_bag_message.hpp→topic_metadata.hpp→bag_metadata.hpp→storage_options.hpp - 插件加载:
storage_factory_impl.hpp→plugin_description.xml(sqlite3) - 写路径:
writer.hpp→sequential_writer.cpp→message_cache.hpp - 读路径:
reader.hpp→sequential_reader.cpp - 压缩:
reader_writer_factory.cpp→sequential_compression_writer.cpp - 录制:
record_options.hpp→recorder.cpp(subscribe_topic/record) - 回放:
play_options.hpp→player.cpp(play/play_messages_from_queue) - CLI:
ros2bag/setup.py→verb/record.py→rosbag2_py/_transport.cpp - 高级:
bag_rewrite.cpp→reindexer.cpp
16. 小结
rosbag2 目录实现 ROS 2 可插拔日志栈:底层 storage 插件 持久化序列化字节;rosbag2_cpp 提供顺序读写、缓存与 metadata;rosbag2_transport 连接 rclcpp 图(GenericSubscription/Publisher)与时间控制;rosbag2_compression 与 ros2bag/rosbag2_py 完成压缩与用户接口。
理解任意 bug 或性能问题,通常从 Recorder::subscribe_topic → Writer::write → storage 插件(录制)或 Player::play → Reader::read_next → GenericPublisher(回放)两条主链入手,并核对 metadata.yaml 与 QoS/serialization 配置 是否一致。
正在加载留言…