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

WebSocket异步更新列表跨进程传递计算实现方案咨询

方案可行性判断

直接将主进程中持续更新的普通Python列表传入独立子进程完全不可行。
Python多进程默认做内存隔离:

  • 用fork模式启动子进程时,子进程仅能拿到fork瞬间主进程列表的快照,后续主进程对列表的更新不会同步到子进程
  • 用Windows/macOS默认的spawn模式启动时,子进程会重新初始化对象,根本拿不到主进程里更新的列表值
    常见的跨进程传值方案比如multiprocessing.Queue、管道、Manager代理列表要么存在序列化/反序列化开销,要么会出现数据堆积,延迟完全达不到25ms级更新的要求,不适合这个场景。
最优实现路径

用共享内存做跨进程数据同步,全程零拷贝,延迟在微秒级,完全适配当前更新频率,不需要改动你现有的WebSocket数据接收核心逻辑。

核心注意事项

  • 共享内存读写必须加进程锁,避免读到半更新的脏数据
  • 锁的持有时间要尽可能短:更新时仅做赋值操作,读取时仅做数据拷贝,所有耗时计算全部放到锁释放后执行,避免阻塞主进程的WebSocket数据接收
  • 多进程代码必须放在if __name__ == "__main__"入口下执行,避免跨平台启动异常

可直接落地的代码实现

代码里已经修正了你原有代码中returnLists方法缺失self参数的问题:

import asyncio
import time
import multiprocessing
from multiprocessing import shared_memory, Lock
import numpy as np

# 全局配置和你原有逻辑对齐
ARRAY_LENGTH = 5
DATA_TYPE = "d"  # 双精度浮点数,可根据实际存储类型调整

def calculation_process(shm_x_name, shm_y_name, proc_lock):
    """独立进程中运行的计算任务"""
    # 绑定主进程创建的共享内存块
    shm_x = shared_memory.SharedMemory(name=shm_x_name)
    shm_y = shared_memory.SharedMemory(name=shm_y_name)
    # 将共享内存映射为可直接操作的数组
    list_x = np.ndarray((ARRAY_LENGTH,), dtype=DATA_TYPE, buffer=shm_x.buf)
    list_y = np.ndarray((ARRAY_LENGTH,), dtype=DATA_TYPE, buffer=shm_y.buf)

    try:
        while True:
            # 加锁拷贝完整的一帧数据,锁释放后再做计算
            with proc_lock:
                calc_x = list_x.copy()
                calc_y = list_y.copy()
            
            # 以下替换为实际计算逻辑
            print(f"计算进程读取到数据 | x: {calc_x}, y: {calc_y}")
            time.sleep(0.025)  # 按更新频率匹配读取节奏即可
    finally:
        # 进程退出前关闭共享内存连接
        shm_x.close()
        shm_y.close()

class DataStore:
    def __init__(self, shm_x, shm_y, store_lock):
        self.lock = store_lock
        # 主进程侧映射共享内存
        self.list_x = np.ndarray((ARRAY_LENGTH,), dtype=DATA_TYPE, buffer=shm_x.buf)
        self.list_y = np.ndarray((ARRAY_LENGTH,), dtype=DATA_TYPE, buffer=shm_y.buf)
        # 初始化默认值和原有逻辑一致
        self.list_x[:] = [0, 0, 0, 0, 0]
        self.list_y[:] = [0, 0, 0, 0, 0]
    
    async def updateList(self, var1, var2, idx):
        # 加锁更新单条数据,执行完立即释放锁
        with self.lock:
            self.list_x[idx] = var1
            self.list_y[idx] = var2

    async def returnLists(self):
        with self.lock:
            return self.list_x.copy(), self.list_y.copy()

async def dataCallback(data, receipt_timestamp, store):
    if data.name == 'B':
        await store.updateList(data.info[0], data.info[1], 0)
    elif data.name == 'C':
        await store.updateList(data.info[0], data.info[1], 1)
    elif data.name == 'D':
        await store.updateList(data.info[0], data.info[1], 2)
    elif data.name == 'E':
        await store.updateList(data.info[0], data.info[1], 3)
    elif data.name == 'F':
        await store.updateList(data.info[0], data.info[1], 4)

async def main():
    # 主进程初始化共享内存块
    elem_size = np.dtype(DATA_TYPE).itemsize
    shm_x = shared_memory.SharedMemory(create=True, size=ARRAY_LENGTH * elem_size)
    shm_y = shared_memory.SharedMemory(create=True, size=ARRAY_LENGTH * elem_size)
    proc_lock = Lock()

    # 启动计算独立进程
    calc_proc = multiprocessing.Process(
        target=calculation_process,
        args=(shm_x.name, shm_y.name, proc_lock),
        daemon=True
    )
    calc_proc.start()

    # 初始化数据存储实例
    store = DataStore(shm_x, shm_y, proc_lock)

    # --- 以下为模拟WebSocket数据接收,替换为你实际的WebSocket启动逻辑即可 ---
    async def mock_ws_stream():
        msg_idx = 0
        name_list = ['B', 'C', 'D', 'E', 'F']
        while True:
            mock_data = type('Data', (), {
                "name": name_list[msg_idx % 5],
                "info": [msg_idx * 1.0, msg_idx * 2.0]
            })()
            await dataCallback(mock_data, time.time(), store)
            msg_idx += 1
            await asyncio.sleep(0.005)  # 模拟5路数据每5ms上报1路,25ms完成一轮全量更新
    # --- 模拟逻辑结束 ---

    try:
        await mock_ws_stream()
    finally:
        # 退出前清理所有资源
        calc_proc.terminate()
        calc_proc.join()
        shm_x.close()
        shm_x.unlink()
        shm_y.close()
        shm_y.unlink()

if __name__ == "__main__":
    asyncio.run(main())

替代方案说明

如果不想引入numpy依赖,可以直接用multiprocessing.Array实现共享内存,底层逻辑完全一致,只是数组操作语法稍有区别,性能没有差异。不要尝试用第三方消息队列、RPC框架做这个场景的数据同步,额外的网络/IO开销会带来不必要的延迟。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:25:05