You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

ROS2话题配置仅保留最新QoS仍留存多条旧消息问题排查

ROS2多节点状态同步异常问题

问题现象

实现支持创建、编辑目标点的规划器时,共部署9个功能节点,各节点负责修改目标点的一类属性,节点列表如下:

/add_target  
/change_comment  
/change_target_index  
/clear_state  
/remove_target  
/rename_target  
/set_target  
/toggle_select_target  
/toggle_visible

所有节点继承统一的StateNode基类,通过共享全局状态实现数据一致,设计逻辑为:

  • 节点收到对应服务调用(如/planner/rename_target)后,从本地存储的状态中定位目标点完成修改,将新状态发布到/planner/state话题
  • 所有节点均订阅/planner/state话题,将收到的消息更新为本地状态,保证全节点状态一致
  • 话题配置QoS为保留最新1条消息,但跨节点调用服务后,使用如下命令查看话题内容:
ros2 topic echo --qos-history keep_last --qos-depth 1 --qos-durability transient_local --qos-reliability reliable /planner/state

会收到多条顺序随机的消息,各节点本地存储状态看似一致,但话题中存在浮动旧消息,不符合QoS仅保留最新消息的预期。

复现步骤

  1. 连续两次调用添加目标点服务:
ros2 service call /planner/add_target mtms_interfaces/srv/AddTarget "{target: {position:{x: 0.0,y: 0.0,z: 0.0}, orientation: {alpha: 0.0,beta: 0.0,gamma: 0.0}}}"

此时话题echo输出正常。
2. 调用重命名目标点服务:

ros2 service call /planner/rename_target mtms_interfaces/srv/RenameTarget "{name: 'Target-0', new_name: 'example'}"

此时话题echo输出2条消息:一条为未修改的旧状态,一条为修改后的新状态。

相关代码实现

StateNode基类

class StateNode(Node):

    def __init__(self, name):
        super().__init__(name)

        # Persist the latest sample.
        qos = QoSProfile(
            depth=1,
            durability=DurabilityPolicy.TRANSIENT_LOCAL,
            history=HistoryPolicy.KEEP_LAST,
            reliability=ReliabilityPolicy.RELIABLE
        )

        self._state_publisher = self.create_publisher(
            PlannerState,
            "/planner/state",
            qos
        )
        self._state_subscriber = self.create_subscription(
            PlannerState,
            '/planner/state',
            self.state_updated,
            10
        )
        self._state = None

    def state_updated(self, msg):
        self._state = msg

RenameTargetNode实现

class RenameTargetNode(StateNode):

    def __init__(self):
        super().__init__('rename_target')
        self.create_service(RenameTarget, '/planner/rename_target', self.rename_target_callback)

    def rename_target_callback(self, request, response):

        state = self._state
        if state is None:
            response.success = False
            return response

        self.get_logger().info('Renaming {} to {}'.format(request.name, request.new_name))
        
        i = 0
        for target in state.targets:

            # Name already exists
            if target.name == request.new_name: 
                response.success = False
                return response
            
            # Save index of target in case new_name is unique
            if target.name == request.name:
                i = state.targets.index(target) 
        
        state.targets[i].name = request.new_name

        self._state_publisher.publish(state)

        response.success = True
        return response

AddTargetNode实现

class AddTargetNode(StateNode):

    def __init__(self):
        super().__init__('add_target')
        self.create_service(AddTarget, '/planner/add_target', self.add_target_callback)

    def first_available_target_name(self):
        if self._state is None:
            return "Target-0"

        target_names = [target.name for target in self._state.targets]
        idx = 0
        while True:
            target_name = "Target-{}".format(idx)
            if target_name not in target_names:
                break
            idx += 1
        return target_name

    def create_new_target(self, pose):
        target = Target()

        target.name = self.first_available_target_name()
        target.type = "Target"
        target.comment = ""
        target.selected = False
        target.target = False  # XXX: Misnomer
        target.pose = pose

        target.intensity = 100.0
        target.iti = 100.0

        return target

    def add_target_callback(self, request, response):
        self.get_logger().info('Incoming request')

        target = self.create_new_target(
            pose=request.target  # XXX: Misnomer
        )

        if self._state is None:
            msg = PlannerState()
            msg.targets = [
                target
            ]
        else:
            msg = self._state
            msg.targets.append(target)

        self._state_publisher.publish(msg)

        response.success = True
        return response

运行环境

  • Ubuntu 20.04,内核版本5.14.0-1042-oem,x86_64架构
  • 所有ROS2节点运行于基于osrf/ros:galactic-desktop镜像创建的Docker容器中

问题根因

问题由两个实现和认知偏差直接导致:

  1. 对ROS2 QoS机制的认知错误
    KEEP_LAST depth=1的队列是每个Publisher实例单独维护的,不是整个话题全局共享一份队列。现在9个功能节点每个都创建了/planner/state的Publisher,等于存在9个独立的发送队列,每个队列各自存1条对应Publisher最后发送的消息。
    配置TRANSIENT_LOCAL持久化策略后,只要有新订阅者(比如执行ros2 topic echo的客户端)连上来,DDS会要求所有在线Publisher把自己队列里存的最后一条消息推给新订阅者。当某个节点处理完服务、发布新状态后,剩下8个节点还没来得及通过订阅回调更新本地状态,它们的Publisher队列里存的还是旧状态消息,新订阅者就会同时收到新消息和多个节点推的旧消息,看起来就是多条顺序混乱的内容。

  2. Python消息对象的引用误用
    ROS2 Python接口里,订阅回调收到的msg是中间件复用的内存对象,直接赋值给self._state只是传了引用,没有做拷贝。后面业务逻辑直接改self._state的属性(比如改target名称、往targets列表追加元素),本质是直接改订阅拿到的原消息对象,既会污染本地状态,还可能因为内存复用出现数据竞争,导致发出去的消息内容不对。另外节点发布完新状态后,没有立刻更新本地存的状态和Publisher队列内容,要等自己发的消息绕一圈经订阅回调回来才更新,这段时间里本地Publisher队列里存的还是旧消息。

修复方案

  • 优先调整架构:撤掉所有功能节点上的/planner/state发布器,全局只留1个独立的状态管理节点当这个话题的唯一发布者。其他功能节点只保留订阅权限,要改状态的时候,把修改请求发给状态管理节点,由它统一做校验、更新状态、发布最新值,从根上避免多发布者各自缓存旧消息的问题。
  • 如果暂时不想改多发布者的架构,必须把引用的坑填上:
    • 订阅回调收到消息后,做深拷贝再存到self._state,别直接存原消息的引用
    • 改状态的时候,基于本地拷贝的状态生成新的状态对象,改完再发布
    • 发布新状态后,立刻把本地self._state更新成刚发布的新状态,别等订阅回调触发了再更新,避免本地发布器缓存旧消息

内容的提问来源于stack exchange,提问作者Alqio

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 20:21:29