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
相关产品推荐
相关产品推荐

