数据聚合应用多进程共享预分配内存的DataFrame/Numpy数组方案咨询
解决方案:跨进程共享动态增长的OHLC DataFrame
一、使用multiprocessing.shared_memory预分配固定内存块
完全可行,步骤如下:
- 预计算内存需求:每个OHLC结构最终固定1000行,按你代码中的12列
float64计算,单个数组的内存大小为1000 * 12 * 8 = 96000字节(float64占8字节)。 - 创建共享内存并映射numpy数组:
主进程中,对每个键(如EURUSD-Deriv-S1)创建固定大小的SharedMemory实例,将numpy数组直接映射到该共享内存,替代直接用np.zeros分配内存:
把共享内存的名称、数组形状、dtype等元数据保存到import numpy as np from multiprocessing import shared_memory key_name = "EURUSD-Deriv-S1" rows = 1000 cols = 12 dtype = np.float64 nbytes = rows * cols * dtype.itemsize # 创建共享内存 shm = shared_memory.SharedMemory(create=True, size=nbytes, name=f"OHLC_{key_name}") # 将numpy数组映射到共享内存 np_array = np.ndarray((rows, cols), dtype=dtype, buffer=shm.buf) # 初始化数组(可选,填充0) np_array.fill(0)multiprocessing.Manager创建的字典中,供子进程读取。 - 子进程读取共享内存:
# 子进程中通过元数据重新映射共享内存 shm = shared_memory.SharedMemory(name=f"OHLC_{key_name}") np_array = np.ndarray((rows, cols), dtype=dtype, buffer=shm.buf) # 转成只读DataFrame使用 df = pd.DataFrame(np_array, columns=["time", "open", "high", "low", "close", ...])
二、Ray框架方案(更简洁)
Ray的对象存储天然支持跨进程共享数据,无需手动管理共享内存:
- 初始化Ray:
import ray ray.init() - 主进程存储DataFrame:
动态更新DataFrame直到达到1000行,每次更新后用ray.put()将最新的DataFrame存入对象存储:# 假设df_dict是你的DataFrame字典 shared_refs = {} for key, df in df_dict.items(): shared_refs[key] = ray.put(df) # 存储只读引用 - 子进程读取数据:
Ray会自动处理内存共享,避免数据拷贝,适配动态增长的场景。@ray.remote def process_ohlc(key): df = ray.get(shared_refs[key]) # 获取只读DataFrame # 处理逻辑... return result # 调用子进程任务 results = ray.get([process_ohlc.remote(key) for key in shared_refs.keys()])
三、修复你之前的numpy数组丢失问题
你的代码存在变量错误:
- 第一个循环的
key_name是循环变量,结束后保留最后一个值;第二个循环中你仍用key_name索引numpyframes_dict,导致所有修改只作用于最后一个键的数组,其他键的数组看似“丢失”。 - 修正代码:将第二个循环中的
numpyframes_dict[key_name][-1,:] = new_numpy_line改为numpyframes_dict[dfname][-1,:] = new_numpy_line,dfname才是当前循环对应的键。
内容的提问来源于stack exchange,提问作者HermDP
相关产品推荐
相关产品推荐

