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

如何在ROS2中实现多话题消息同步?解决订阅节点同步问题

ROS2话题消息同步与配对实现方案

问题分析

你当前的Combine节点存在两个核心问题:

  • 消息配对逻辑错误:收到新的单话题消息时,会直接复用另一话题的历史消息,导致输出的消息并非严格配对
  • 无超时/停止检测机制:话题停止发布后仍会重复使用旧消息,无法实现"输出None"或"等待重新发布"的需求

实现方案

通过消息队列维护未配对消息+定时超时检测的方式,完全匹配你需要的三个场景。以下是修改后的完整代码:

import rclpy
from rclpy.node import Node
from std_msgs.msg import String
from collections import deque

class Combine(Node):
    def __init__(self):
        super().__init__("combine")
        # 订阅两个话题
        self.subs_hello_ = self.create_subscription(
            String, '/hello', self.callback_msg_hello, 10)
        self.subs_world_ = self.create_subscription(
            String, '/world', self.callback_msg_world, 10)
        
        # 维护两个消息队列,保存未配对的消息
        self.hello_queue = deque()
        self.world_queue = deque()
        
        # 创建定时器,每1秒检查一次消息队列(可根据需求调整间隔)
        self.timer = self.create_timer(1.0, self.check_timeout_and_output)

    def callback_msg_hello(self, msg):
        # 收到hello消息,先尝试和world队列的消息配对
        if self.world_queue:
            world_msg = self.world_queue.popleft()
            self.get_logger().info(f"Hello: {msg.data}, World: {world_msg.data}")
        else:
            # 没有可配对的world消息,存入队列
            self.hello_queue.append(msg)

    def callback_msg_world(self, msg):
        # 收到world消息,先尝试和hello队列的消息配对
        if self.hello_queue:
            hello_msg = self.hello_queue.popleft()
            self.get_logger().info(f"Hello: {hello_msg.data}, World: {msg.data}")
        else:
            # 没有可配对的hello消息,存入队列
            self.world_queue.append(msg)

    def check_timeout_and_output(self):
        # 场景2:仅一个队列有消息,输出带None的内容
        if self.hello_queue and not self.world_queue:
            # 取出所有未配对的hello消息(这里取最新的,也可根据需求取全部)
            hello_msg = self.hello_queue.pop()
            self.get_logger().info(f"Hello: {hello_msg.data}, World: None")
            # 清空队列,避免重复输出
            self.hello_queue.clear()
        elif self.world_queue and not self.hello_queue:
            world_msg = self.world_queue.pop()
            self.get_logger().info(f"Hello: None, World: {world_msg.data}")
            self.world_queue.clear()
        # 场景3:两个队列都为空,不做任何操作,等待新消息
        else:
            pass

def main(args=None):
    rclpy.init(args=args)
    combine_node = Combine()
    rclpy.spin(combine_node)
    combine_node.destroy_node()
    rclpy.shutdown()

if __name__ == '__main__':
    main()

关键逻辑说明

  1. 消息配对逻辑

    • 每次收到新消息时,优先尝试从另一话题的队列中取出最早的消息进行配对输出,保证消息的顺序性和配对准确性
    • 无配对消息时,将当前消息存入对应队列等待后续配对
  2. 超时检测与场景处理

    • 通过定时器定期检查队列状态:
      • 仅一个队列有消息时,取出最新消息并输出带None的内容,同时清空队列避免重复输出
      • 两个队列都为空时,不执行任何操作,等待新消息发布后再处理
  3. 边界情况处理

    • 双话题正常发布时,消息会一一配对输出,符合场景1需求
    • 单话题停止发布时,定时器触发后输出对应消息+None,符合场景2需求
    • 双话题都停止时,队列清空后无输出,直到两者重新发布新消息才会继续配对,符合场景3需求

内容的提问来源于stack exchange,提问作者Mr. x

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 16:48:13