【问题标题】:python multiprocessing - method not invoked with expected argumentspython multiprocessing - 未使用预期参数调用的方法
【发布时间】: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


    【解决方案1】:

    我找到了解决办法。问题是调用函数应该有 str 类型的参数。如果你传递一些像字典这样的对象,它就不能正常工作。

    【讨论】:

      【解决方案2】:

      p.map 并不像您想象的那样工作。

      当您调用p.map(function,data) 时,如果数据是一个数组(如您的情况),那么池将在data 的每个元素上运行function:

      def saveOutputs(data):
           print(data)
      
      out={1001:"dummy", 1002:"foo", 1003:"bar"}
      
      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)
              p.map(saveOutputs,data)
      

      会给你:

      [1001, 'dummy']
      1001
      dummy
      [None, None]
      [1002, 'foo']
      1002
      foo
      [None, None]
      [1003, 'bar']
      1003
      bar
      [None, None]
      

      对于前一对数据,两次调用function,每个调用都带有一对各自的元素。

      【讨论】:

      • 谢谢琼。我期望 Pool 在我的数据的每个元素上运行。如果我以您的数据为例,我面临的问题是我的函数没有以完整的值 1001 传递。1001 的各个数字被拆分为 1,0,0,1,然后用于进一步处理。这破坏了我的系统。
      • @Techie 不应该,但我们不知道您的字典包含什么...
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-07
      • 1970-01-01
      • 2012-11-03
      相关资源
      最近更新 更多