libstatistics_collector 源码详细分析
工作区路径:/home/cp/work2/ros2Learn/ros2_humble/src/ros-tooling/libstatistics_collector
版本:1.3.4,许可证 Apache 2.0,Quality Level 1。
libstatistics_collector 是 ROS 2 轻量级统计聚合库:提供在线滑动窗口统计(均值/最大/最小/标准差)、通用 Collector 框架、话题级 message_age / message_period 采集器,以及将结果封装为 statistics_msgs/MetricsMessage 的工具。主要消费者是 rclcpp 的 Topic Statistics 功能。
1. 仓库结构
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23
| libstatistics_collector/ ├── include/libstatistics_collector/ │ ├── collector/ │ │ ├── collector.hpp # 抽象采集器基类 │ │ ├── metric_details_interface.hpp # 指标名/单位接口 │ │ └── generate_statistics_message.hpp # → MetricsMessage │ ├── moving_average_statistics/ │ │ ├── moving_average.hpp # Welford 在线统计 │ │ └── types.hpp # StatisticData │ ├── topic_statistics_collector/ │ │ ├── topic_statistics_collector.hpp # 话题采集器模板接口 │ │ ├── received_message_age.hpp # 消息年龄 │ │ ├── received_message_period.hpp # 消息周期 │ │ └── constants.hpp # 指标名/参数名常量 │ └── visibility_control.hpp ├── src/libstatistics_collector/ │ ├── collector/collector.cpp │ ├── collector/generate_statistics_message.cpp │ └── moving_average_statistics/{moving_average,types}.cpp ├── test/ # gtest + benchmark ├── CMakeLists.txt / package.xml ├── README.md / QUALITY_DECLARATION.md └── Doxyfile
|
构建产物: 单一共享库 liblibstatistics_collector.so(话题采集器为 header-only 模板)。
2. 在 ROS 2 栈中的位置
| 层级 |
职责 |
| MovingAverageStatistics |
纯数学:O(1) 内存聚合 |
| Collector |
生命周期 + 聚合器封装 |
| TopicStatisticsCollector |
从 ROS 消息提取度量值 |
| rclcpp::SubscriptionTopicStatistics |
定时发布、窗口管理 |
3. 核心模块一:MovingAverageStatistics
3.1 算法
使用 Welford 在线算法 计算总体标准差,无需存储全部样本:
1 2 3 4 5 6 7 8 9 10 11
| void MovingAverageStatistics::AddMeasurement(const double item) { if (!std::isnan(item)) { count_++; const double previous_average = average_; average_ = previous_average + (item - previous_average) / count_; min_ = std::min(min_, item); max_ = std::max(max_, item); sum_of_square_diff_from_mean_ += (item - previous_average) * (item - average_); } }
|
标准差:sqrt(sum_of_square_diff_from_mean_ / count_)(总体标准差,非样本标准差)。
3.2 StatisticData
1 2 3 4 5 6 7
| struct StatisticData { double average = NaN; double min = NaN; double max = NaN; double standard_deviation = NaN; uint64_t sample_count = 0; };
|
无样本时返回 NaN;NaN 输入会被 丢弃。
3.3 “Moving Average” 含义
名称易误解:并非固定长度的滑动窗口,而是 当前采集窗口 内的在线统计;窗口结束需调用 Reset() / ClearCurrentMeasurements() 清零开始新窗口。
4. 核心模块二:Collector
4.1 类层次
1 2 3 4 5 6 7 8 9
| MetricDetailsInterface ├── GetMetricName() └── GetMetricUnit() ↑ Collector (abstract) ├── AcceptData(double) ├── GetStatisticsResults() ├── Start() / Stop() └── SetupStart() / SetupStop() [纯虚]
|
4.2 生命周期
1 2 3 4 5 6 7 8
| bool Collector::Start() { lock → if already started return false started_ = true → SetupStart() } bool Collector::Stop() { lock → started_ = false → SetupStop() ClearCurrentMeasurements() // 在锁外调用 Reset }
|
| 方法 |
行为 |
AcceptData(m) |
直接写入 MovingAverageStatistics(不检查 started_) |
GetStatisticsResults() |
返回当前窗口统计 |
ClearCurrentMeasurements() |
collected_data_.Reset() |
Stop() |
停止并清空测量 |
4.3 线程安全
Collector 的 mutex_ 保护 started_ 与 Start/Stop
MovingAverageStatistics 自带 mutex_ 保护统计数据
AcceptData 不加 Collector 锁,仅依赖内部聚合器锁
5. 核心模块三:GenerateStatisticMessage
将 StatisticData 转为 ROS 消息:
1 2 3 4 5 6
| MetricsMessage GenerateStatisticMessage(node_name, metric_name, unit, window_start, window_stop, data) { // 填充 5 个 StatisticDataPoint: // AVERAGE, MAXIMUM, MINIMUM, SAMPLE_COUNT, STDDEV }
|
| 字段 |
来源 |
measurement_source_name |
节点名 |
metrics_source |
如 "message_age" |
unit |
如 "ms" |
window_start/stop |
采集窗口时间 |
statistics[] |
avg/min/max/count/stddev |
6. 核心模块四:Topic Statistics Collectors
6.1 接口
1 2 3 4 5
| template<typename T> class TopicStatisticsCollector : public collector::Collector { virtual void OnMessageReceived(const T & msg, const rcl_time_point_value_t now_nanoseconds) = 0; };
|
6.2 ReceivedMessageAgeCollector
度量: 接收时刻 − 消息 header.stamp(毫秒)
1 2 3 4 5 6 7
| void OnMessageReceived(const T & msg, rcl_time_point_value_t now_ns) { auto [has_header, stamp_ns] = TimeStamp<T>::value(msg); if (has_header && stamp_ns && now_ns) { age_ms = (now_ns - stamp_ns) in milliseconds; AcceptData(age_ms); } }
|
编译期检测 header:
HasHeader<M>:是否存在 M.header.stamp 且类型为 builtin_interfaces/msg/Time
- 无 header 的消息:静默跳过(不记录)
| 属性 |
值 |
GetMetricName() |
"message_age" |
GetMetricUnit() |
"ms" |
6.3 ReceivedMessagePeriodCollector
度量: 相邻两次 OnMessageReceived 调用的时间间隔(毫秒)
1 2 3 4 5 6 7 8 9 10
| void OnMessageReceived(..., now_ns) { lock(mutex_); if (time_last_ == kUninitializedTime) time_last_ = now_ns; // 首条消息只初始化,不产生样本 else { period_ms = (now_ns - time_last_) in ms; time_last_ = now_ns; AcceptData(period_ms); } }
|
| 属性 |
值 |
GetMetricName() |
"message_period" |
GetMetricUnit() |
"ms" |
| 线程安全 |
自有 mutex_(与 Collector 锁独立) |
6.4 常量
1 2 3 4 5
| kMsgAgeStatName = "message_age" kMsgPeriodStatName = "message_period" kMillisecondUnitName = "ms" kCollectStatsTopicNameParam = "collect_topic_name" kPublishStatsTopicNameParam = "publish_topic_name"
|
7. rclcpp 集成(下游)
rclcpp::topic_statistics::SubscriptionTopicStatistics 封装完整工作流:
1 2 3 4 5 6 7 8 9 10 11 12 13
| bring_up(): ReceivedMessageAge → Start() ReceivedMessagePeriod → Start()
handle_message(msg, now): for collector: collector->OnMessageReceived(msg, now.ns())
publish_message_and_reset_measurements(): // 定时器触发,默认 1s for collector: stats = GetStatisticsResults() ClearCurrentMeasurements() publish GenerateStatisticMessage(...) window_start = window_end
|
默认发布话题:/statistics
启用方式:创建 subscription 时 SubscriptionOptions.enable_topic_statistics = true(及节点参数配置)。
8. 依赖关系
1 2 3 4 5
| libstatistics_collector ├── rcl # rcl_time_point_value_t, RCL_S_TO_NS ├── rcpputils # thread_safety_annotations ├── statistics_msgs # MetricsMessage, StatisticDataType └── builtin_interfaces # Time(header 检测)
|
不依赖 rclcpp——保持库层轻量,由 rclcpp 在上层集成。
9. 测试
| 测试 |
覆盖 |
test_moving_average_statistics |
Welford 正确性、NaN、Reset |
test_collector |
Start/Stop、AcceptData |
test_received_message_age |
有/无 header 消息 |
test_received_message_period |
周期间隔、首条跳过 |
benchmark_iterative |
AddMeasurement 性能 |
测试用 libstatistics_collector_test_msgs(DummyMessage / DummyCustomHeaderMessage)验证 header 检测。
10. 设计特点与局限
| 特点 |
说明 |
| O(1) 内存 |
不存原始样本,适合高频 topic |
| 模板 header-only 话题采集 |
易扩展新 metric |
| 与 ROS 消息格式对齐 |
直接生成 MetricsMessage |
| QL1 声明 |
有 QUALITY_DECLARATION |
| 局限 |
说明 |
| message_age 需 header.stamp |
无 header 类型无法统计 |
| age 依赖时钟一致 |
now 与 header 须同源且单调 |
| period 首条不计入 |
每个窗口第一条只建立基准 |
| AcceptData 不检查 started_ |
理论上 Stop 后仍可写入 |
| window 时间用 system_clock |
rclcpp 层非 RCL 时钟 |
| 名称 “moving average” |
实为可重置窗口,非固定长度滑动 |
| 仅 subscriber 侧 metric |
无 publisher 延迟等内置采集器 |
11. 扩展自定义 Collector 示例
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
| class MyMetricCollector : public libstatistics_collector::collector::Collector { protected: bool SetupStart() override { return true; } bool SetupStop() override { return true; } public: std::string GetMetricName() const override { return "my_metric"; } std::string GetMetricUnit() const override { return "ms"; }
void OnEvent(double value_ms) { AcceptData(value_ms); } };
MyMetricCollector c; c.Start(); c.OnEvent(12.3); auto stats = c.GetStatisticsResults(); auto msg = GenerateStatisticMessage("my_node", c.GetMetricName(), c.GetMetricUnit(), t0, t1, stats);
|
话题级可继承 TopicStatisticsCollector<T> 并实现 OnMessageReceived。
12. 推荐阅读顺序
moving_average.hpp + moving_average.cpp — 核心算法
collector.hpp + collector.cpp — 生命周期模式
generate_statistics_message.cpp — ROS 消息映射
received_message_age.hpp / received_message_period.hpp — 话题 metric
rclcpp/.../subscription_topic_statistics.hpp — 端到端集成
test_moving_average_statistics.cpp — 算法边界条件
13. 小结
libstatistics_collector 是 ROS 2 Topic Statistics 的底层数学与抽象层:MovingAverageStatistics 用 Welford 算法 O(1) 聚合样本;Collector 提供 Start/Stop 与指标元数据;两个内置话题采集器计算 message_age 与 message_period;GenerateStatisticMessage 输出标准 MetricsMessage。rclcpp 在其上实现定时发布,使订阅者可向 /statistics 报告通信质量指标。
如需,我可以把本文写入 ros2doc/ros-tooling/libstatistics_collector源码详细分析.md,或继续分析 rclcpp 中 enable_topic_statistics 的完整启用路径与参数。
正在加载留言…