如何在Python中实现跨文件传递持续更新的实时数据
跨文件流式数据传递方案
针对实时行情接收+计算的场景,按实现复杂度、稳定性从高到低给三个可落地方案,按需选择即可。
方案1:同进程异步队列传递(首推,零依赖、延迟最低)
不需要把两个文件拆成独立进程,直接调整入口逻辑,用Python标准库asyncio.Queue做数据缓冲即可,既不会阻塞websocket接收,也不会出现数据丢失,完全满足实时行情的传输需求。
代码修改示例
首先调整B.py,把websocket接收逻辑封装成可复用的协程,不再作为独立入口:
import asyncio import json import time import websockets async def binance_ws_consumer(data_queue: asyncio.Queue, symbol: str = "ethusdt"): url = f'wss://stream.binance.com:9443/stream?streams={symbol}@miniTicker' async with websockets.connect(url) as client: while True: raw_data = json.loads(await client.recv())['data'] event_time = time.localtime(raw_data['E'] // 1000) event_time_str = f"{event_time.tm_hour}:{event_time.tm_min}:{event_time.tm_sec}" latest_price = float(raw_data['c']) # 数据放入队列,立即返回不阻塞接收 await data_queue.put((event_time_str, latest_price))
然后把A.py作为唯一启动入口,启动时初始化队列、拉起B的接收协程作为后台任务,再循环从队列取数据做计算:
import asyncio from B import binance_ws_consumer async def run_strategy(): # 初始化队列,maxsize可根据计算速度调整,避免内存溢出 data_queue = asyncio.Queue(maxsize=100) # 后台启动websocket接收任务 asyncio.create_task(binance_ws_consumer(data_queue)) while True: # 阻塞等待新数据,不占CPU资源 event_time, price = await data_queue.get() # 以下写你的计算逻辑即可 print(f"[接收数据] 时间:{event_time} 最新价:{price}") data_queue.task_done() if __name__ == "__main__": asyncio.run(run_strategy())
方案优势
- 无跨进程通信开销,无序列化/反序列化成本,传输延迟在微秒级
- 队列自带缓冲,就算计算逻辑有毫秒级耗时,也不会阻塞websocket接收导致断连或者丢包
- 全用Python标准库能力,不需要安装额外第三方依赖,问题排查成本极低
- 可灵活配置队列满后的处理策略,比如丢弃最旧数据优先保证实时性,非常适配行情场景
方案2:多进程队列传递(适配CPU密集型计算场景)
如果你的计算逻辑是CPU密集型(比如大量指标计算、AI推理、复杂回测),会卡住异步事件循环影响websocket连接稳定性,就把两个逻辑拆成独立进程,用multiprocessing.Queue做跨进程数据传递:
- A作为主进程,启动时拉起B的子进程运行websocket接收逻辑
- B进程收到数据后直接放入跨进程队列
- A主进程循环从队列读取数据做计算,计算逻辑卡主也不会影响B进程的接收
注意:进程间内存是隔离的,不要尝试用全局变量、普通函数调用的方式跨进程传数据,必须用Queue、Pipe这类专门的IPC机制。
方案3:本地消息中间件传递(适配多模块消费场景)
如果后续你需要多个独立策略、多个模块同时消费这一份行情数据,可以用本地轻量消息通道做中转:B收到数据后统一推送到消息通道,所有需要数据的模块自行订阅消费即可,模块间完全解耦。
避坑提醒:不要用本地文件读写轮询、HTTP轮询这类方式传高频流式数据,延迟高还容易出现读写冲突,稳定性极差。
内容的提问来源于stack exchange,提问作者mohammadreza
相关产品推荐
相关产品推荐

