返回专辑
·Johan·15 分钟阅读

ContentFilteredTopic:把过滤下推到中间件还是应用层

对比订阅端丢弃、应用层条件判断与 DDS 内容过滤的 CPU/带宽边界,说明表达式能力与发现时序限制;用高频率噪声话题与延迟尖峰验收过滤位置选择。

ContentFilteredTopic:把过滤下推到中间件还是应用层

1. 过滤位置是系统设计,不是回调里的 if

一条 /detections/raw 从前端感知节点以 200 Hz 发布,每帧包含几十个候选框、置信度和调试特征。规划只关心 class_id == 1score > 0.7 的目标,监控只关心低置信度样本,记录器则希望拿全量。最直接的实现是在每个订阅回调开头写:

cpp
if (msg->score <= min_score_ || msg->class_id != target_class_) {
  return;
}

功能看起来正确,系统负载却没有降下来:网卡仍在收全量样本,DDS reader 仍要处理缓存与可靠性状态,ROS 2 类型支持仍可能完成反序列化,executor 仍被唤醒。回调里的 return 只是让算法少算一步,不等于让通信链路少搬一次数据。

ContentFilteredTopic 的价值在于把「这条样本是否值得投递」变成中间件可见的订阅条件。它不是 QoS 的替代品,也不是通用查询引擎,而是 DDS 对 topic 的一个过滤视图:订阅端声明表达式与参数,RMW 尝试把条件交给 DDS。如果实现支持,过滤可以发生在 reader 收样本前,甚至由 writer 按每个匹配订阅者裁剪发送集合。位置越靠近发布端,越能省带宽;位置越靠近应用,越能使用业务上下文。工程判断就在这两端之间。

2. 三种常见过滤层

第一层是应用层条件判断。消息已经进入用户回调,字段可随便读,能访问参数、状态机、TF、缓存地图,表达能力最强。缺点也最清楚:网络、DDS 缓存、序列化、executor 调度基本都已经付费。它适合低频话题、复杂业务条件、一次性调试,以及过滤比例不高的场景。

第二层是订阅端丢弃。有些项目会在独立的过滤节点或订阅封装里统一丢弃样本,例如先订阅全量 /raw,再发布 /filtered。这能把下游算法从噪声中隔离出来,便于复用与监控,却没有降低发布者到过滤节点之间的带宽。若过滤节点和消费者不在同一进程,还会多一次发布订阅。它更像拓扑整理,不是通信优化。

第三层是 DDS ContentFilteredTopic。订阅者在创建时给出类似 SQL 的谓词,例如 score > %0 AND class_id = %1,并提供参数。理想路径是 DDS 在传输或 reader 队列之前判断样本是否匹配。若发布者在远端,且 DDS vendor 支持 writer-side filtering,非匹配样本甚至不会上网线;若只支持 reader-side filtering,网络带宽仍消耗,但用户回调、ROS 队列和部分反序列化成本会下降。不要把「用了 CFT」直接等同于「带宽下降」,它必须用抓包或发送端计数验证。

3. 成本模型:先算哪一段可能被省掉

把一条消息的代价拆成四段:发布端序列化 、链路传输 、订阅端接收与反序列化 、业务处理 。原始频率为 ,过滤通过率为 。应用层判断的近似成本是:

如果 DDS 能在 reader 入队前过滤,但仍收到全量网络数据,则变成:

其中 是表达式求值成本。若 vendor 能把过滤下推到 writer 侧,远端订阅者的链路成本也按 缩小:

这组公式不是为了得到精确百分比,而是提醒测试该看哪里。大消息、低通过率、跨机器传输时,CFT 才可能明显改善带宽与尾延迟;小消息、同机进程、过滤比例接近一时,表达式求值和发现复杂度可能比收益更醒目。把过滤下推之前,先记录 topic 的平均字节数、峰值频率、通过率和消费者数量,否则很容易把一个可读的 if 换成不可解释的中间件行为。

4. rclcpp 里声明过滤条件

ROS 2 的 C++ 入口通常放在 SubscriptionOptions。表达式是字符串,参数是字符串数组,DDS 在底层按类型解释字段。下面只展示关键形状,字段名应来自消息定义,而不是 C++ 成员函数或业务别名:

cpp
rclcpp::SubscriptionOptions options;
options.content_filter_options.filter_expression =
  "score > %0 AND class_id = %1";
options.content_filter_options.expression_parameters = {"0.70", "1"};

sub_ = this->create_subscription<vision_msgs::msg::Detection2D>(
  "/detections/raw",
  rclcpp::SensorDataQoS(),
  std::bind(&DetectorConsumer::on_detection, this, std::placeholders::_1),
  options);

运行时调整阈值时,优先使用订阅对象提供的内容过滤更新接口,而不是销毁再创建订阅。销毁重建会触发重新发现,短时间内可能出现匹配空窗;参数更新则更接近「同一个 reader 改条件」。但这里仍要做失败处理:有的 RMW 会返回不支持,有的表达式能创建却不能按预期下推,有的只在部分类型上工作。初始化日志应明确打印过滤表达式、参数和 RMW 实现,诊断里也要暴露当前条件。

不要把 ros2 topic echo --filter 当成 CFT 验收。CLI 过滤常用于人工观察,条件可能在 Python 层对已收到的消息求值;它证明「这个谓词写法能筛出样本」,不证明 DDS 已经减少网络或 reader 负载。真正验收要写一个使用 SubscriptionOptions 的小节点,或在目标节点上打开同一条配置路径。

5. 表达式能力:像 SQL,但不是数据库

DDS 内容过滤语言通常是 SQL-92 条件表达式的子集:比较、逻辑与或、括号、字符串匹配、参数占位符是常见能力。它面向单条样本的字段判断,不面向跨样本聚合,也不能访问 ROS 参数、TF、当前时间、外部地图或上一帧状态。score > %0 合适,distance_to_goal(pose) < 2.0 不合适;class_id IN (1, 2) 可能取决于实现,moving_average(score, 5) 则不属于这类机制。

字段路径也要保守。基础标量字段最稳,嵌套结构、数组下标、字符串函数、枚举映射和 bounded sequence 的支持会随着 RMW 与 DDS vendor 变化。对于 std_msgs/msg/Float32 这种只有一个字段的消息,表达式通常写 data > %0;对于自定义消息,建议把高频过滤所需字段设计成顶层标量,例如:

cpp
# Message: NoiseSample.msg
# 顶层字段便于 DDS 过滤;复杂数组留给应用层处理。
std_msgs/Header header
uint8 source_id
float32 noise_rms
float32 latency_ms
float32[] spectrum

如果过滤条件需要读 spectrum[17] 或根据 header.frame_id 做复杂模式匹配,说明消息边界可能需要重做:把路由字段提升成顶层元数据,或者拆分 topic。CFT 适合低成本判定「这条样本属于哪类消费者」,不适合替代算法分类器。表达式越像业务规则引擎,跨厂商风险越高。

这里还有一个经常被忽略的契约:过滤字段一旦进入 CFT,就不再只是内部实现细节,而是 topic 的路由协议。发布者随手把 score 从归一化置信度改成未校准 logit,应用层回调也许能通过参数补偿,DDS 表达式却仍按旧阈值丢样本。自定义消息应把可过滤字段写进注释或接口文档,说明单位、取值范围、缺省值和版本迁移策略。若字段语义仍在快速变化,不要急着下推;先把 topic 拆成稳定元数据和实验载荷,或在发布端显式产生面向不同消费者的派生话题。

过滤字段还要避免承担过多含义。source_id == 3 可以是硬件通道,severity >= 2 可以是诊断级别;但把多个条件编码进一个 uint32 flags,再要求表达式做位运算,通常会撞上实现差异。宁可多放几个顶层标量,也不要让 CFT 依赖晦涩的位布局。消息字节数增加一点,换来表达式可读、跨语言一致、压测可解释,往往是值得的。

6. 发现与匹配时序:过滤不是瞬时全局状态

DDS 发现是异步的。订阅者创建后,reader 的 QoS、类型信息和内容过滤条件会通过发现协议传播给 publisher;publisher 看到这个 reader、确认兼容并建立匹配需要时间。CFT 因此有两个容易误判的窗口。

第一个窗口在启动阶段。订阅节点刚起来时,发布者可能还不知道过滤条件,或者只知道有 reader 尚未拿到完整参数。不同实现会选择先不发、先按未过滤发送、或在 reader 侧兜底过滤。对于 volatile 数据,这通常只表现为前几帧统计波动;对于 transient local 或需要历史样本的场景,历史缓存如何套用过滤条件要单独验证。不要把「节点启动后一秒内收到一条不该收到的样本」简单归咎于业务代码,先看发现时间线。

第二个窗口在动态更新参数时。阈值从 0.7 改到 0.9,应用层变量立即生效,但 DDS 过滤条件需要更新到 reader,并可能传播给 writer。期间的样本属于旧条件还是新条件,不应承载安全语义。若过滤决定会影响制动、互锁或故障恢复,必须在应用层再做一次最终校验,把 CFT 当减负手段而不是安全边界。

工程上可以把过滤条件版本写进诊断:每次更新递增 filter_generation,回调里记录最后一次收到样本的 generation 和时间。验收时同时抓 ros2 topic info -v、节点日志和网络计数,才能区分「未匹配」「匹配但未下推」「表达式错误」这三类问题。

同机组合也需要单独验证。ROS 2 进程内通信、组件容器、loaned message 和共享内存传输都会改变样本经过的路径;某个 DDS 层面的优化在跨主机 UDP 上有效,不代表在同一 component container 里仍然有效。若发布者和订阅者被组合进同一进程,最便宜的过滤位置可能反而是发布端提前分流,或者在算法入口保留普通 if。判断标准不是 API 名字,而是实际路径上哪段成本最大。

启动顺序测试要覆盖「publisher 先起」「subscriber 先起」「过滤参数运行时改变」「网络短暂断开后重连」四种情况。每种情况都记录第一条合法样本的到达时间、是否出现不匹配样本、是否发生回调空窗。CFT 引入的复杂度主要藏在这些边缘时序里,平稳运行十分钟没有异常,并不能证明重连后一秒也安全。

7. QoS 与 CFT 的边界

QoS 决定能否匹配、缓存多深、丢旧还是阻塞;CFT 决定匹配后哪类样本值得投递。二者经常一起出现,却不能互相替代。BestEffort + KeepLast(1) 能压低滞后,但不能避免低价值样本占满链路;CFT 能减少低价值样本,但若 Reliable 在弱网下等待 ACK,剩余样本仍可能形成尾延迟。过滤通过率下降不代表 deadline 自动更安全,尤其是发布者仍以高频产生样本时。

对可靠传输还要问一个细节:非匹配样本是否参与 reader 的可靠性状态?理想情况下,writer 知道某 reader 的过滤条件后,不需要为该 reader 跟踪非匹配样本的确认;如果过滤只能在 reader 侧做,网络和可靠性开销就不会消失。这个行为不是 ROS 2 API 能一眼看出的,必须看 DDS 实现说明或用压测观察发送端吞吐、重传、history 占用。

生命周期节点也有相似边界。inactive 状态下是否保留订阅、何时更新过滤参数、activate 前是否已经完成发现,都会影响第一批样本。对高频传感链路,建议把订阅创建、过滤参数设置、激活顺序写成 launch test:先启动 publisher,再启动 consumer;反过来再跑一次。两种启动顺序都稳定,才算过滤位置选择可落地。

8. 厂商差异:把支持矩阵写进风险清单

ROS 2 的 API 给了统一入口,实际执行落在 RMW 和 DDS vendor 上。Fast DDS、Cyclone DDS、Connext 等实现对内容过滤的支持深度、表达式子集、writer-side 下推、动态更新和 introspection 类型处理并不完全相同。更麻烦的是,同一 vendor 的不同版本也可能修复或改变过滤路径。

因此 CFT 不适合悄悄引入公共库后全仓默认开启。更稳的做法是按话题建一张支持矩阵:消息类型、表达式、RMW、DDS 版本、是否跨主机省带宽、是否支持运行时更新、失败时的降级路径。矩阵里的结论必须来自本项目构建和运行环境,不来自博客复制。若生产镜像允许切换 RMW_IMPLEMENTATION,CI 至少要覆盖主用实现;备选实现要么跑同样压测,要么明确标注「只保证功能,不保证下推收益」。

降级路径也要提前设计。CFT 创建失败时,可以退回应用层判断并打出高优先级告警;表达式不支持某个字段时,可以改用过滤节点;跨主机带宽没有下降时,可以拆 topic 或把路由字段前移到发布端。最坏的方案是静默退回全量订阅,让团队以为已经省下带宽,现场却在高负载时暴露尖峰。

日志要能支撑这种降级判断。节点启动时至少打印 RMW 实现、过滤表达式、参数、创建结果和降级状态;运行中定期输出收到样本数、通过样本数、回调丢弃数和最近一次参数更新时间。若 CFT 生效后应用层最终校验仍丢掉大量样本,说明表达式与业务条件不一致;若 CFT 通过率很低但网卡吞吐不变,说明下推位置不够靠前。把这些指标接入 diagnostic_updater 或 Prometheus,比事后翻 DDS trace 容易得多。

对录包也要谨慎。ros2 bag record /noise/raw 记录的是 recorder 自己订阅到的视图;如果 recorder 也用了过滤条件,bag 就不再代表全量现场。排查性能时建议同时录两路:一条全量小窗口采样,用于还原发布端事实;一条过滤后长时间记录,用于观察消费者实际负载。回放时若消费者仍带 CFT,要确认 bag 播放器作为 publisher 是否支持对应发现行为,否则 replay 与现场可能出现不同的过滤位置。

9. 高频噪声话题的验收方法

验收不应只看平均 CPU。构造一个高频噪声话题,让绝大多数样本对目标消费者无用,同时保留少量突发样本模拟真实事件。例如发布 /noise/raw,字段包含 source_idnoise_rmslatency_ms 和一段频谱数组;消费者只关心 source_id == 3 AND noise_rms > 0.8。发布频率从 100 Hz、500 Hz、1000 Hz 逐档上升,消息体大小也从几十字节扩到数 KB。

python
import random
import rclpy
from rclpy.node import Node
from example_interfaces.msg import Float32MultiArray

class NoisePublisher(Node):
    def __init__(self):
        super().__init__("noise_publisher")
        self.pub = self.create_publisher(Float32MultiArray, "/noise/raw", 10)
        self.timer = self.create_timer(0.001, self.tick)

    def tick(self):
        msg = Float32MultiArray()
        source_id = random.randint(0, 7)
        rms = random.random()
        latency = random.choice([2.0, 3.0, 5.0, 80.0])
        msg.data = [float(source_id), rms, latency] + [random.random() for _ in range(256)]
        self.pub.publish(msg)

def main():
    rclpy.init()
    rclpy.spin(NoisePublisher())

同一套 publisher 下跑三组 consumer:回调内 if、过滤节点转发、CFT 订阅。每组至少记录五类数据:publisher 进程 CPU、consumer 进程 CPU、跨主机网卡吞吐、回调实际频率、端到端延迟 P50/P99/P999。若 CFT 只让 consumer CPU 下降而网卡不变,说明过滤发生在 reader 侧或更晚;若 publisher CPU 上升而网卡下降,说明 writer 正在为不同 reader 求值,这可能仍然划算,但要看总 CPU 是否转移到关键节点。

采样窗口要足够长,且必须包含负载扰动。可以在同一台机器上给消费者加一个低优先级 CPU 压力线程,在网络侧用限速或丢包模拟无线链路,再观察三组方案的 P999 是否分叉。应用层 if 往往在轻载下表现很好,因为业务处理被跳过后平均 CPU 降了;一旦 executor 被全量唤醒、reader 队列开始积压,尾延迟会先坏。CFT 的收益也可能只在跨主机时出现,所以本机压测通过后还要把 publisher 放到另一台机器或容器网络里重跑。

通过率统计不能只看消费者收到多少条。要在 publisher 侧用同一套谓词离线计算理论匹配数,再和 consumer 侧收到数比较。若理论匹配数为每秒 20 条,CFT consumer 收到 18 条,应用层 consumer 收到 20 条,就要查 reliability、history depth、deadline 和表达式精度,而不是直接宣布过滤省了资源。性能优化不能以丢掉合法样本为代价。

延迟尖峰比平均值更重要。过滤的目标通常不是让 demo 更省电,而是在系统接近饱和时避免低价值样本挤占关键路径。压测时故意加入 80 ms 的 latency_ms 尖峰,并让消费者只接收尖峰或只排除尖峰,观察 executor 队列深度、callback 间隔和业务线程调度。若应用层判断的 P999 仍被全量消息唤醒拖高,而 CFT 组 P999 收敛,才说明下推位置对实时性有价值。

最终报告建议画三条曲线:横轴是发布频率,纵轴分别是网卡吞吐、consumer CPU 和 P999 延迟。应用层过滤通常吞吐曲线随频率线性上升,CPU 在某个点后被唤醒成本拖住;过滤节点会把下游曲线压低,但过滤节点自身成为新瓶颈;真正下推成功的 CFT 应该让目标订阅者的吞吐和回调频率接近理论通过率。只要其中一条曲线不符合预期,就继续定位,不要只凭「回调次数少了」做结论。

还要为「没有收益」设定退出条件。若在目标硬件上,CFT 与应用层过滤的 P999 差异落在测量噪声内,跨主机吞吐也没有下降,就保留普通条件判断,并把原因写进话题设计记录。中间件特性不是越多越成熟;能用简单代码解释清楚的低频路径,通常不值得引入发现时序和厂商矩阵。CFT 应留给瓶颈明确、通过率低、压测曲线已经证明下推有效的链路。这样的退出标准能保护代码可读性,也能避免性能讨论变成对某个中间件特性的信仰。

10. 选择规则:先保证正确,再争取下推

实际项目可以按四个问题决策。第一,过滤条件是否只依赖单条消息的稳定字段?若需要 TF、地图、状态机或跨帧统计,留在应用层。第二,通过率是否足够低、消息是否足够大、消费者是否跨主机?若三者都不是,CFT 的收益可能不值得引入差异风险。第三,当前 RMW 是否证明能在目标路径下降低网络或 reader 成本?若没有证据,就把它当功能过滤而不是性能优化。第四,过滤错误的后果是什么?若错误会漏掉安全相关样本,应用层必须保留最终判定。

一个稳妥落地顺序是:先在回调里写清晰条件并加统计,得到真实通过率;再把路由字段整理成顶层标量,避免表达式碰复杂结构;然后为单个高频 topic 增加 CFT,并保留应用层断言;最后用高频噪声与延迟尖峰压测决定是否推广。推广时每个 topic 都要有「表达式、参数来源、降级行为、验收数据」四项记录。

ContentFilteredTopic 解决的是通信系统里的选择位置问题。它让中间件有机会少搬无用样本,却要求团队接受表达式子集、发现时序和厂商实现差异。把它用在高频、低通过率、跨主机的大消息上,并用 CPU、带宽和尾延迟三条曲线验收,收益会很清楚;把它当成回调 if 的漂亮替代品,则很容易得到一个更难调试、收益不明的过滤层。

← 全部文章

johan's blog