【发布时间】:2020-06-06 15:06:17
【问题描述】:
我有一个字典对象,其中包含一个输出,其中键为“id”,值为 pandas 数据框。字典的大小为 9。我需要将 pandas 数据帧的输出保存在 HDFS 上每个 id 的单个文件中。考虑到将每个文件写入 13 分钟 * 9 = 107 分钟所需的时间,我试图将其并行化,以便每个文件的写入并行发生。
作为这个用例的一部分,我正在尝试使用多处理,如下所示 -
def saveOutputs(data):
print(data[0])
#logic to write data in file
with Pool(processes = 9) as p:
for k, v in out.items(): #out is a dict which i need to persist in file
data = [k,v]
print(data[0])
p.map(saveOutputs,data)
我看到的是,如果我的 id(dict 中的键) 是 1001 ,当 saveOutputs 作为 print 的一部分在 saveOutputs 中被调用时,它会将值打印为 1 而不是 1001 而在我的 Pool 块中调用 saveOutputs 之前,打印语句正在打印1001.
我对这种行为不是很清楚,也不确定不正确的地方缺少什么。 寻找一些输入。
谢谢。
【问题讨论】:
标签: python-3.x pandas pyspark python-multiprocessing