多进程向Asyncio传递Cryptofeed OHLC数据问题求助
解决多进程中cryptofeed数据无法共享给主进程asyncio任务的问题
问题描述
尝试用cryptofeed模块获取OHLC数据,将数据流放在独立多进程中,想把数据存到全局变量后,用独立的asyncio实例访问该变量。但多进程无法共享全局数据,async函数close()输出的是空pandas DataFrame。
原代码:
from cryptofeed import FeedHandler from cryptofeed.backends.aggregate import OHLCV from cryptofeed.defines import TRADES from cryptofeed.exchanges import BinanceFutures import pandas as pd from multiprocessing import Process from concurrent.futures import ProcessPoolExecutor import asyncio data1 = pd.DataFrame() # Create an empty DataFrame queue = multiprocessing.Queue() async def ohlcv(data): global data1 # Convert data to a Pandas DataFrame df = pd.DataFrame.from_dict(data, orient='index') # Reset the index df.reset_index(inplace=True) df.index = [pd.Timestamp.now()] data1 = data1.append(df) queue.put('nd') # Append the rows of df to data async def close(data): while True: print(data) await asyncio.sleep(15) def main1(): f = FeedHandler() f.add_feed(BinanceFutures(symbols=['BTC-USDT-PERP'], channels=[TRADES], callbacks={TRADES: OHLCV(ohlcv, window=10)})) f.run() if __name__ == '__main__': p = Process(target=main1) p.start() asyncio.run(close(data1))
问题原因
多进程的内存空间是相互独立的,子进程里修改全局变量data1只会影响子进程自己的内存副本,主进程中的data1不会被更新,所以close()输出空DataFrame。你创建的Queue没有被用来传递实际数据,只是传了无意义的字符串,无法实现数据共享。
解决方案
利用multiprocessing.Queue在子进程和主进程之间传递处理好的DataFrame片段,主进程的asyncio任务从队列中取出数据并维护共享的DataFrame。修改后的代码如下:
from cryptofeed import FeedHandler from cryptofeed.backends.aggregate import OHLCV from cryptofeed.defines import TRADES from cryptofeed.exchanges import BinanceFutures import pandas as pd from multiprocessing import Process, Queue import asyncio def ohlcv_callback(data, queue): # 处理OHLC数据生成DataFrame片段 df = pd.DataFrame.from_dict(data, orient='index') df.reset_index(inplace=True) df.index = [pd.Timestamp.now()] # 将处理好的DataFrame放入队列传递给主进程 queue.put(df) def main1(queue): f = FeedHandler() # 将队列传入回调函数,实现进程间数据传递 f.add_feed(BinanceFutures( symbols=['BTC-USDT-PERP'], channels=[TRADES], callbacks={TRADES: OHLCV(lambda d: ohlcv_callback(d, queue), window=10)} )) f.run() async def close(queue): shared_df = pd.DataFrame() while True: # 循环取出队列中所有数据,非阻塞避免任务卡住 while not queue.empty(): df_chunk = queue.get() # 使用concat替代已弃用的append方法拼接数据 shared_df = pd.concat([shared_df, df_chunk]) print("当前累积的OHLC数据:") print(shared_df) await asyncio.sleep(15) if __name__ == '__main__': queue = Queue() # 将队列作为参数传给子进程 p = Process(target=main1, args=(queue,)) p.start() asyncio.run(close(queue))
修改说明
- 将
Queue作为参数传递给子进程和回调函数,让子进程能向主进程传递实际数据 - 子进程回调函数不再修改全局变量,而是把处理好的DataFrame片段放入队列
- 主进程的
close函数维护本地的shared_df,循环从队列取数据并拼接,实现数据共享 - 用
pd.concat替代已被弃用的DataFrame.append方法,保证代码兼容性
内容的提问来源于stack exchange,提问作者Cosmin George
相关产品推荐
相关产品推荐

