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

多进程下实时传感器数据访问延迟问题及优化方案咨询

多进程传感器实时数据延迟问题解决方案

问题描述

我是多进程技术新手,在父脚本中启动子进程作为传感器,该子进程持续将实时数据发送至管道,供父进程按需获取。但我发现存在数据延迟问题,获取到的并非传感器实时数据而是历史数据。添加time.sleep(0.1)后问题略有改善,但精度仍未达标。请问还有哪些处理实时数据的可行方案?


子脚本函数

我添加了time.sleep(0.1),问题略有改善但精度仍不足

def _send_data(self):
        time.sleep(1)
        offsets = self.calculate_offset()
        
        try:
            while 1:
               
                if self._conn:
                    zeroed_data = self.zero_data(offsets,self.get_data())
                    status = {
                        "Status": self._status,
                        "Timestamp": sR.millis(),
                        "Temperature": self._temperature,
                        }
                    status.update(zeroed_data)
                self._conn.send(status)
                time.sleep(0.1)

        except KeyboardInterrupt:
            # ctrl-C abort handling
            if self.conn:
                self.conn.close() 
            print('stopped')



def main_bota(child_conn,port):
    print('bota started')
    bota_sensor_1 = None
   
    try:
        bota_sensor_1 = BotaSerialSensor(port,child_conn);
        bota_sensor_1.run()
       
    except BotaSerialSensorError as expt:
        print('bota failed: ' + expt.message)
    
    except Exception as e:
        print(f"An unexpected error occured: {e}")
    
    finally:
        if bota_sensor_1 is not None:
            bota_sensor_1.close()
        sys.exit(1)

父脚本函数

def robot_command(selected_tool, path):
    """send commands to robot according to the path.

    Args:
        selected_tool (np.array): selected tool matrix (4x4)
        path (list): List of transformation matrices defining the robots path.
    """

    worldToToolStart = sR.getWorldToTool(s, BUFFER_SIZE, selected_tool)
    totNumbPos = len(path)
    totNumb = 0
   
    for i in range(totNumbPos):
        totNumb += 1
        toolToNewTarget = path[i]
        worldToNewTarget = np.matmul(worldToToolStart, toolToNewTarget)
        worldToFlangeCommand = np.matmul(
            worldToNewTarget, sR.homInv(selectedTool))

        # Formulate the command based on robot mode
        if robotMode == 'SmartServo':
            poseStringForRobot = sR.make_cmd(
                "MovePTPHomRowWiseMinChangeRT", worldToFlangeCommand)

        else:
            poseStringForRobot = sR.make_cmd(
                "MovePTPHomRowWiseMinChange", worldToFlangeCommand)


        # Send command to robot
        sR.command(s, BUFFER_SIZE, poseStringForRobot)
      
        # get sensor data
        sensordata = parent_conn.recv()
       
        # Check status of Bota
        #status_bota(sensordata['Status'])
       
        # get robot data dict and combine with sensor data dict
        robotdata = fetch_robot_data()
        combined_data = {**sensordata, **robotdata}
        
        #append combined_data to global data
        data.append(combined_data)
        print(i)
        print(robotdata["timestamp"])
        print(sensordata["Timestamp"])
        # Wait until position is reached (can be adjusted...)
        sR.delay(delayBetweenPoses)

传感器连接与断开函数

def connect_bota(port):
    global p,parent_conn
    parent_conn, child_conn = Pipe()
    p = Process(target=main_bota, args = (child_conn,port))
    p.start()
    print('process bota started')


def disconnect_bota():
    p.terminate()
    p.join()
    print('process bota stopped')



def main():
    config()
    name = generate_experiment_name()
    print(name)
    display_config()
    connect_to_robot()
    connect_bota(port='COM3')
    
    # generate selected path
    path = generate_path(selectedPath_str)

    # (go to supportRobot file to adapt initialization parameters (velocities, damping, ...))
    sR.initRobot(s, BUFFER_SIZE, robotMode)
    print('initRobot')

    # Delay for 3s
    delay(2)

    # Send command to robot
    print('send command')
    robot_command(selectedTool, path)

    # End of operations
    print('STOP in 2s')
    delay(2)
    disconnect_bota()
    s.close()
    print("STOPPED")

可行解决方案

1. 清空管道缓存,只取最新数据

管道会缓存子进程发送的所有数据,父进程每次recv()只会拿到最早的未读数据,导致历史数据堆积。可以在获取数据前先清空缓存,只保留最新的一条:

# 父进程获取传感器数据时修改为:
while parent_conn.poll():
    # 读取并丢弃所有堆积的历史数据
    parent_conn.recv()
# 最后一次读取就是最新数据
sensordata = parent_conn.recv()

这种方式不需要修改子进程逻辑,快速解决数据堆积问题。

2. 改为"请求-响应"模式,按需采集

让子进程不再主动推送数据,而是等待父进程的请求信号,收到请求后立即采集并发送最新数据,从根源避免堆积:

  • 子进程修改_send_data函数:
def _send_data(self):
    time.sleep(1)
    offsets = self.calculate_offset()
    try:
        while 1:
            # 等待父进程的请求信号
            request = self._conn.recv()
            if request == "GET_LATEST_DATA" and self._conn:
                zeroed_data = self.zero_data(offsets, self.get_data())
                status = {
                    "Status": self._status,
                    "Timestamp": sR.millis(),
                    "Temperature": self._temperature,
                }
                status.update(zeroed_data)
                self._conn.send(status)
    except KeyboardInterrupt:
        if self.conn:
            self.conn.close() 
        print('stopped')
  • 父进程获取数据时先发送请求:
# 发送请求信号
parent_conn.send("GET_LATEST_DATA")
# 接收最新采集的数据
sensordata = parent_conn.recv()

3. 使用共享内存替代管道

管道基于消息队列,存在传输延迟和缓存问题;共享内存是进程间直接读写内存区域,延迟更低,适合高实时性场景:

  • 修改连接函数,创建共享字典:
from multiprocessing import Manager

def connect_bota(port):
    global p, shared_sensor_data
    manager = Manager()
    shared_sensor_data = manager.dict()
    # 子进程参数改为共享字典
    p = Process(target=main_bota, args=(shared_sensor_data, port))
    p.start()
    print('process bota started')
  • 子进程修改_send_data函数,直接更新共享内存:
def _send_data(self, shared_dict):
    time.sleep(1)
    offsets = self.calculate_offset()
    try:
        while 1:
            zeroed_data = self.zero_data(offsets, self.get_data())
            status = {
                "Status": self._status,
                "Timestamp": sR.millis(),
                "Temperature": self._temperature,
            }
            status.update(zeroed_data)
            # 实时更新共享字典
            for k, v in status.items():
                shared_dict[k] = v
            # 可根据传感器性能调整间隔,比管道场景更小
            time.sleep(0.01)
    except KeyboardInterrupt:
        print('stopped')
  • 父进程直接读取共享内存中的最新数据:
# 直接复制共享字典中的实时数据
sensordata = dict(shared_sensor_data)

4. 优化子进程数据发送频率

如果父进程处理每个路径点的时间小于子进程的time.sleep(0.1)间隔,会导致管道数据堆积。可以调小子进程的sleep间隔,或者根据父进程的处理速度动态调整,甚至移除sleep(需确保传感器硬件能承受高频采集)。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 16:29:50