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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 22:54:17