在最近的这段时间里, 我发布了一个宏大的愿望, 这个愿望的意图是想去构建一套专门用于对企业金融展开研究的框架。
后来我把相关的进度情况拉出来查看以后, 才发现该框架的版本已经更新到了三点八这一阶段, 于是我就在这个过程中发现了在三八点八的这个新版本里面新出现的一个模块, 也就是这个句末使用的句号符号本身。
我随手编写了一个测试程序, 这个程序会生成一个大小为240MB的文件, 该文件的具体后缀名可能因系统差异有所不同,接着把这个文件转换成了numpy数组格式。在这个numpy数组里, 它是包含有整型类型数据列和字符串类型数据列的。
之所以选择使用numpy数据类型来做这步处理, 其核心目的是为了能够完整保存各个字段的数据类型信息, 只有这样操作, 才能在后续的子进程环境中, 从左共享内存区域正确无误地读取相关数据内容。
而我在子进程内部进行实际运作时, 使用的是非常简单的计算方法, 也就是简单地直接调用numpy相关的工具来进行计算过程。第一种办法是采用共享内存的方式, 另一种办法是把numpy这个模块当作参数直接传递给子进程。下图所展示的是测试代码的输出的画面。
通过观察可以发现, 利用共享内存去实施第一种途径的时候, 它所消耗掉的内存数量仅仅是可以完全不需要去担心的微不足道的极少部分而已, 并且它能够在2秒这么短的时间以内就把战斗给彻底结束掉。
而另外那种依靠传递参数的手段呢, 它的结果是完全不一样的, 因为它不仅实实在在地占用了高达1.8GB的内存容量, 而且速度缓慢得让人无法忍受, 最终竟然耗费了长达200多秒的漫长岁月。
这很可能跟本人所使用的测试机器有关, 那台测试机是一台二零一七年的十二英寸产品, 配备的是四核的i5处理器与八克内存, 该产品目前已经停产了, 不过这也从一个侧面证明了当数据量达到非常庞大的程度时, 尽量去避免那些没有必要的复制操作和信息传递行为, 确实是一种非常有效的做法。
测试代码如下:
frommultiprocessing.shared_memoryimportSharedMemory
frommultiprocessing.managersimportSharedMemoryManager
fromconcurrent.futuresimportProcessPoolExecutor, as_completed
frommultiprocessingimportcurrent_process, cpu_count
fromdatetimeimportdatetime
importnumpy as np
importpandas as pd
importtracemalloc
importtime
defwork_with_shared_memory(shm_name, shape, dtype):
print(f'With SharedMemory: {current_process()=}')
#Locate the shared memory by its name
shm =SharedMemory(shm_name)
#Create the np.recarray from the buffer of the shared memory
np_array = np.recarray(shape=shape, dtype=dtype, buf=shm.buf)
returnnp.nansum(np_array.val)
defwork_no_shared_memory(np_array: np.recarray):
print(f'No SharedMemory: {current_process()=}')
#Without shared memory, the np_array is copied into the child process
returnnp.nansum(np_array.val)
if__name__=="__main__":
#Make a large data frame with date, float and character columns
a =[
(datetime.today(), 1,'string'),
(datetime.today(), np.nan,'abc'),
] * 5000000
df = pd.DataFrame(a, columns=['date','val','character_col'])
#Convert into numpy recarray to preserve the dtypes
np_array = df.to_records(index=False)
deldf
shape, dtype =np_array.shape, np_array.dtype
print(f"np_array's size={np_array.nbytes/1e6}MB")
#With shared memory
#Start tracking memory usage
tracemalloc.start()
start_time =time.time()
with SharedMemoryManager() as smm:
#Create a shared memory of size np_arry.nbytes
shm =smm.SharedMemory(np_array.nbytes)
#Create a np.recarray using the buffer of shm
shm_np_array = np.recarray(shape=shape, dtype=dtype, buf=shm.buf)
#Copy the data into the shared memory
np.copyto(shm_np_array, np_array)
#Spawn some processes to do some work
with ProcessPoolExecutor(cpu_count()) as exe:
fs =[exe.submit(work_with_shared_memory, shm.name, shape, dtype)
for_inrange(cpu_count())]
for_inas_completed(fs):
pass
#Check memory usage
current, peak =tracemalloc.get_traced_memory()
print(f"Current memory usage {current/1e6}MB; Peak: {peak/1e6}MB")
print(f'Time elapsed: {time.time()-start_time:.2f}s')
tracemalloc.stop()
#Without shared memory
tracemalloc.start()
start_time =time.time()
with ProcessPoolExecutor(cpu_count()) as exe:
fs =[exe.submit(work_no_shared_memory, np_array)
for_inrange(cpu_count())]
for_inas_completed(fs):
pass
#Check memory usage
current, peak =tracemalloc.get_traced_memory()
print(f"Current memory usage {current/1e6}MB; Peak: {peak/1e6}MB")
print(f'Time elapsed: {time.time()-start_time:.2f}s')
tracemalloc.stop()
需要特别做出说明和强调的是, 使用numpy库的时候, 其对应的数据类型一定不能是那种特殊情况的状态, 否则当多个子进程去访问共享内存资源的时候, 程序运行环境必定会出现严重的错误或错误现象, 但是如果在主进程中去直接访问操作共享内存的话, 则完全不会出现任何问题和麻烦状况。
为了进行补充更新, 需要明确指出上面的测试代码里或者说是ory函数的内部实现中是绝对不能够进行解引用操作的, 具体的表现就是绝对不能去调用像print()这样的函数来进行操作, 如果非要这样做的话, 那么程序就会出错。
然而使用.val和.date这两个属性来操作则是完全没有问题的, 之所以不会出现问题的根本原因是因为我们知道这两个属性对应的数据的种类也就是我们常说的dtype并不是那种会导致错误的类型。
相对而言, 通过其他方式得到的数据的dtype在我们的这个特定的场景和这个特定的代码里是属于那种会引发问题的类型的。
解决这个问题, 办法其实也非常简, 在经过无数次踩坑以后, 只需要在小括号内部指定数据类型即可。
= df.(index=False,={'': 'S6'})
在这里我们将长度指定为 6, 因此它需要是一个由 6 个字符组成的字符串。如果满足这个条件的话你就可以把 S6 替换成 U6。一旦字符串的长度超过了规定的限度, 那部分超出范围的字符则会被直接处理掉或者是移除掉。
这样做以后, 那个情况就没有了。核心的关键点其实就是要避免使用特定数据类型的东西。
补充一点点
当完全排除了涉及字符串这种类型的可能性之后, 情况看起来似乎是处于正确状态的。
当涉及到字符串类型的时候, 数据类型设置会变成对应的类型, 这样的话就会导致在处理数据的过程中出现错误, 无法得到正确的处置结果。
这是因为
将参数选项设置为索引列不导出的时候。
它将会把属于字符串类型的数据转换成特定的其他类型。
可以先利用事先给定的dType, 把dType先构造好, 然后再。
npv =
将数据框对象中通过特定方法调用所获取的每一项结果, 逐个提取出来并转换为元组格式。
= np.array(npv, dtype=(dType))
其实呢, 这个方法它是构造得更加正确了。
另外,
如果 str 不指定长度, 似乎这个操作也是不正确的。
也就是说, np.dtype应该是。
如果不这么做的话, 就会导致其他的进程从里面读取出来的数据是空的字符串。
在其他的那些进程当中使用。
shm = ('')
数组对象的形状是包含一万个大小的元组, 数据类型设置为rType类型, 缓冲区参数引用shm对象的buf属性。
df = pd.()
即可以重建
PS1:
这篇文章将为你提供有关<如何在<中<使用>>的具体指导, 希望能为你带来帮助。
PS2:
当其他的进程发生退出的时候, 共享内存是会被自动地进行删除处理的。
为了能够让各个进程各自去使用共享内存, 并且彼此之间互不干扰, 那么还得在进程最终退出之前, 把对共享内存的引用给去掉, 从而避免那种出现尴尬事情的过程出来, 比如当进程自行退出的时候, 如果自动将共享内存进行删除的话, 其他进程就没有办法再进行使用了。