如何在Python多进程中用共享内存实现大型DataFrame只读访问
如何用Python 3.8的共享内存实现多进程只读访问大型DataFrame?
我正在用Python 3.8写第一个multiprocessing程序,有一个大型DataFrame需要给所有进程使用,且所有进程仅需只读访问。我知道共享内存是可行方案,但不清楚具体怎么让每个进程通过共享内存关联到这个DataFrame。
我初步写了这段代码:
# 我的DataFrame,我知道创建共享内存对象前需要先获取总字节数 df_bytes = df.memory_usage(index=True).sum() from multiprocessing import shared_memory shm = shared_memory.SharedMemory(create=True, size=df_bytes) memory_name = shm.name
根据我的理解,共享内存对象会被分配一个后续可用的名称。之后在各进程调用的函数里,我尝试这样关联共享内存:
existing_shm = shared_memory.SharedMemory(name=memory_name)
但我不知道怎么把这个共享内存和我的大型DataFrame绑定?我参考过下面的numpy示例,但不想把DataFrame整体转成numpy数组,而且每个进程都重新创建数组的做法在我看来不太合理。
参考的numpy示例:
>>> # 在第一个Python交互shell中 >>> import numpy as np >>> a = np.array([1, 1, 2, 3, 5, 8]) # 从已有numpy数组开始 >>> from multiprocessing import shared_memory >>> shm = shared_memory.SharedMemory(create=True, size=a.nbytes) >>> # 创建基于共享内存的numpy数组 >>> b = np.ndarray(a.shape, dtype=a.dtype, buffer=shm.buf) >>> b[:] = a[:] # 将原始数据复制到共享内存 >>> b array([1, 1, 2, 3, 5, 8]) >>> type(b) <class 'numpy.ndarray'> >>> type(a) <class 'numpy.ndarray'> >>> shm.name # 未指定名称,系统自动分配 'psm_21467_46075' >>> # 在同一shell或同一机器的新Python shell中 >>> import numpy as np >>> from multiprocessing import shared_memory >>> # 关联到已有的共享内存块 >>> existing_shm = shared_memory.SharedMemory(name='psm_21467_46075') >>> # 注意此示例中a的shape为(6,),dtype为np.int64 >>> c = np.ndarray((6,), dtype=np.int64, buffer=existing_shm.buf) >>> c array([1, 1, 2, 3, 5, 8]) >>> c[-1] = 888 >>> c array([ 1, 1, 2, 3, 5, 888])
请问怎么把大型DataFrame存入共享内存,实现多进程只读访问?
解决方案
方法1:基于DataFrame列的共享内存拆分(避免整体转数组)
DataFrame的每一列本质是numpy数组,我们可以为每一列单独创建共享内存块,在子进程中重新组装成DataFrame,既保留DataFrame结构,又避免全量数据拷贝。
主进程代码(初始化共享内存)
import pandas as pd import numpy as np from multiprocessing import shared_memory import multiprocessing as mp # 假设你的大型DataFrame为df df = pd.read_csv("large_dataset.csv") # 存储每一列的共享内存元信息:名称、形状、数据类型、列名 col_shared_info = [] for col in df.columns: series = df[col] arr = series.to_numpy() # 创建共享内存块 shm = shared_memory.SharedMemory(create=True, size=arr.nbytes) # 创建绑定共享内存的数组并写入数据 shared_arr = np.ndarray(arr.shape, dtype=arr.dtype, buffer=shm.buf) shared_arr[:] = arr[:] # 记录元信息 col_shared_info.append({ "name": shm.name, "shape": arr.shape, "dtype": str(arr.dtype), "col_name": col }) # 单独处理索引的共享内存 index_arr = df.index.to_numpy() index_shm = shared_memory.SharedMemory(create=True, size=index_arr.nbytes) shared_index = np.ndarray(index_arr.shape, dtype=index_arr.dtype, buffer=index_shm.buf) shared_index[:] = index_arr[:] index_info = { "name": index_shm.name, "shape": index_arr.shape, "dtype": str(index_arr.dtype) }
子进程函数(读取共享内存并组装DataFrame)
def worker(col_info_list, index_info): import pandas as pd import numpy as np from multiprocessing import shared_memory # 从共享内存重建每一列 cols_dict = {} for info in col_info_list: # 关联到已有的共享内存块 shm = shared_memory.SharedMemory(name=info["name"]) # 从共享内存创建数组 arr = np.ndarray(info["shape"], dtype=np.dtype(info["dtype"]), buffer=shm.buf) cols_dict[info["col_name"]] = arr # 重建索引 index_shm = shared_memory.SharedMemory(name=index_info["name"]) index_arr = np.ndarray(index_info["shape"], dtype=np.dtype(index_info["dtype"]), buffer=index_shm.buf) index = pd.Index(index_arr) # 组装成可只读访问的DataFrame shared_df = pd.DataFrame(cols_dict, index=index) # 此处可进行只读操作,比如统计、筛选 print(shared_df.describe()) # 注意:子进程不要关闭共享内存,由主进程统一释放
启动并管理进程
if __name__ == "__main__": # 创建4个进程(可根据需求调整数量) processes = [] for _ in range(4): p = mp.Process(target=worker, args=(col_shared_info, index_info)) processes.append(p) p.start() # 等待所有进程执行完毕 for p in processes: p.join() # 主进程统一释放共享内存 # 释放列的共享内存 for info in col_shared_info: shm = shared_memory.SharedMemory(name=info["name"]) shm.close() shm.unlink() # 释放索引的共享内存 index_shm = shared_memory.SharedMemory(name=index_info["name"]) index_shm.close() index_shm.unlink()
方法2:使用multiprocessing.Manager(代码简洁,性能略低)
如果对性能要求不极致,可使用Manager创建DataFrame代理,本质是进程间通信而非直接共享内存,但代码更简单:
import pandas as pd import multiprocessing as mp def worker(shared_df): # 只读访问共享DataFrame print(shared_df.head()) if __name__ == "__main__": df = pd.read_csv("large_dataset.csv") with mp.Manager() as manager: # 创建共享的DataFrame代理 shared_df = manager.Namespace() shared_df.df = df # 启动进程 processes = [mp.Process(target=worker, args=(shared_df,)) for _ in range(4)] for p in processes: p.start() for p in processes: p.join()
关键注意事项
- 因为是只读访问,无需处理数据竞争,不用加锁。
- 共享内存的关闭和释放必须由主进程完成,子进程仅关联使用即可。
- 若DataFrame包含
object类型的字符串列,需先转为pandas的StringDtype(1.0+版本支持),否则共享内存无法直接处理指针类型数据;此类场景也可优先选择Manager方案。
内容的提问来源于stack exchange,提问作者mHelpMe
相关产品推荐
相关产品推荐

