【问题标题】:Similar errors in MultiProcessing. Mismatch number of arguments to functionMultiProcessing 中的类似错误。函数的参数数量不匹配
【发布时间】:2016-10-12 07:59:30
【问题描述】:

我找不到更好的方法来描述我所面临的错误,但每次我尝试将多处理实现到循环调用时似乎都会出现此错误。

我使用过 sklearn.externals.joblib 和 multiprocessing.Process,但错误相似但不同。

要应用多处理的原始循环,其中一个迭代在单个线程/进程中执行

for dd in final_col_dates:
    idx1 = final_col_dates.tolist().index(dd)

    dataObj = GetPrevDataByDate(d1, a, dd, self.start_hour_of_day)
    data2 = dataObj.fit()

    dataObj = GetAppointmentControlsSchedule(data2, idx1, d, final_col_dates_mod, dd, self.DC, frgt_typ_filter)
    data3 = dataObj.fit()

    if idx1 > 0:
       data3['APPT_SCHD_ARVL_D_{}'.format(idx1)] = np.nan

    iter += 1

    days_out_vars.append(data3)

为了将上面的代码片段实现为多处理,我创建了一个方法,上面的代码除了 for 循环

使用Joblib,下面是我的代码sn-p。

Parallel(n_jobs=2)(
            delayed(self.ParallelLoopTest)(dd, final_col_dates, d1, a, d, final_col_dates_mod, iter, return_list)
                    for dd in final_col_dates)

变量return_list是在方法ParallelLoopTest中执行的共享变量。它被声明为:

manager = Manager()
return_list = manager.list()

使用上面的代码sn-p,我遇到以下错误:

Process SpawnPoolWorker-3:
Traceback (most recent call last):
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\process.py", line 249, in _bootstrap
  self.run()
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\process.py", line 93, in run
  self._target(*self._args, **self._kwargs)
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\pool.py", line 108, in worker
  task = get()
File "C:\Users\dkanhar\Anaconda3\lib\site-packages\sklearn\externals\joblib\pool.py", line 359, in get
  return recv()
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\connection.py", line 251, in recv
  return ForkingPickler.loads(buf.getbuffer())
TypeError: function takes at most 0 arguments (1 given)

我也尝试了多处理模块来执行上述代码,但仍然遇到类似的错误。以下代码用于使用多处理模块运行:

for dd in final_col_dates:
    # multiprocessing.Pipe(False)
    p = multiprocessing.Process(target=self.ParallelLoopTest, args=(dd, final_col_dates, d1, a, d, final_col_dates_mod, iter, return_list))
    jobs.append(p)
    p.start()

for proc in jobs:
    proc.join()

而且,我面临以下错误回溯:

File "<string>", line 1, in <module>
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\spawn.py", line 106, in spawn_main
   exitcode = _main(fd)
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\spawn.py", line 116, in _main
   self = pickle.load(from_parent)
TypeError: function takes at most 0 arguments (1 given)
Traceback (most recent call last):
File "E:/Projects/Predictive Inbound Cartoon Estimation-MLO/Python/dataprep/DataPrep.py", line 457, in <module>
   print(obj.fit())
File "E:/Projects/Predictive Inbound Cartoon Estimation-MLO/Python/dataprep/DataPrep.py", line 39, in fit
return self.__driver__()
File "E:/Projects/Predictive Inbound Cartoon Estimation-MLO/Python/dataprep/DataPrep.py", line 52, in __driver__
   final = self.process_()
File "E:/Projects/Predictive Inbound Cartoon Estimation-MLO/Python/dataprep/DataPrep.py", line 135, in process_
   sch_dat = self.inline_apply_(all_dates_schd, d1, d2, a)
File "E:/Projects/Predictive Inbound Cartoon Estimation-MLO/Python/dataprep/DataPrep.py", line 297, in inline_apply_
   p.start()
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\process.py", line 105, in start
   self._popen = self._Popen(self)
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\context.py", line 212, in _Popen
   return _default_context.get_context().Process._Popen(process_obj)
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\context.py", line 313, in _Popen
   return Popen(process_obj)
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\popen_spawn_win32.py", line 66, in __init__
   reduction.dump(process_obj, to_child)
File "C:\Users\dkanhar\Anaconda3\lib\multiprocessing\reduction.py", line 59, in dump
   ForkingPickler(file, protocol).dump(obj)
   BrokenPipeError: [Errno 32] Broken pipe

所以,我尝试取消注释 multiprocessing.Pipe(False) 行,认为这可能是因为使用了我禁用的 Pipe,但问题仍然存在,我面临同样的错误。

如果有任何帮助,以下是我的方法 ParallerLoopTest:

def ParallelLoopTest(self, dd, final_col_dates, d1, a, d, final_col_dates_mod, iter, days_out_vars):
    idx1 = final_col_dates.tolist().index(dd)

    dataObj = GetPrevDataByDate(d1, a, dd, self.start_hour_of_day)
    data2 = dataObj.fit()

    dataObj = GetAppointmentControlsSchedule(data2, idx1, d, final_col_dates_mod, dd, self.DC, frgt_typ_filter)
    data3 = dataObj.fit()

    if idx1 > 0:
        data3['APPT_SCHD_ARVL_D_{}'.format(idx1)] = np.nan

    print("Iter ", iter)
    iter += 1

    days_out_vars.append(data3)

我之所以说类似的错误是因为如果您查看两个错误的 Traceback ,它们之间都有类似的错误行:

TypeError:函数最多接受 0 个参数(给定 1 个),而我不知道为什么会发生这种情况。

另外请注意,我之前在其他项目中成功实现了这两个模块,但从未遇到过问题,所以我不知道为什么现在开始出现这个问题,以及这个问题究竟意味着什么。

任何帮助都将不胜感激,因为自 3 天以来我一直在浪费时间进行调试。

谢谢

在最后一个答案后编辑 1

回答后,我尝试了以下这个。 添加装饰器 @staticmethod,移除 self,并使用 DataPrep.ParallelLoopTest(args) 调用方法。

另外,将方法移出 DataPrep 类,并由 ParallelLoopTest(args) 简单地调用,

但在这两种情况下,错误仍然相同。

PS:我尝试在这两种情况下都使用 joblib。 所以,这两种解决方案都不起作用。

新方法定义:

def ParallelLoopTest(dd, final_col_dates, d1, a, d, final_col_dates_mod, iter, days_out_vars, DC, start_hour):
    idx1 = final_col_dates.tolist().index(dd)

    dataObj = GetPrevDataByDate(d1, a, dd, start_hour_of_day)
    data2 = dataObj.fit()

    dataObj = GetAppointmentControlsSchedule(data2, idx1, d, final_col_dates_mod, dd, DC, frgt_typ_filter)
    data3 = dataObj.fit()

    if idx1 > 0:
        data3['APPT_SCHD_ARVL_D_{}'.format(idx1)] = np.nan

    print("Iter ", iter)
    iter += 1

    days_out_vars.append(data3)

编辑 2:

我遇到了错误,因为 Python 无法腌制一些大型数据帧。我的参数/参数中有 2 个 DataFrame,一个大约 20MB,另外一个是 pickle 格式的 200MB。但这应该不是问题吧?我们应该能够通过 Pandas DataFrame。如果我错了,请纠正我。

另外,解决方法是我在使用随机名称调用方法之前将 DataFrame 保存为 csv,传递文件名并读取 csv,但这是一个缓慢的过程,因为它涉及到巨大的 csv 文件。有什么建议吗?

【问题讨论】:

  • 您有重现问题的最少代码吗?
  • 我添加的代码是产生错误的必需代码。你说的最小代码是什么意思你能清楚一点吗?
  • 我无法运行您的代码并得到相同的错误,所以我看不出哪里出错了。例如 [这里][gist.github.com/tomMoral/c75824eea5b3f68fd2148c64d1ee88fa] 是一段代码,用于测试 Process 和 Manager 之间的交互。它可以在您的计算机上运行吗?你能腌制你在Process中使用的所有对象吗?
  • 我收到页面未找到错误。你能验证你发布的链接吗?

标签: python multiprocessing pickle joblib


【解决方案1】:

实际上在这两种情况下都会遇到完全相同的错误,但是当您在一个示例中使用 Pool (joblib) 和在另一个示例中使用 Process 时,您不会在主程序中遇到相同的故障/回溯线程,因为它们不会以相同的方式管理进程失败。
在这两种情况下,您的流程似乎都无法在新的Process 中解开您的子作业。 Pool 会返回 unpickling 错误,而使用 Process 则会失败,因为当子进程死于此 unpickling 错误时,它会关闭主线程用于写入数据的管道,从而导致主进程出错.

我的第一个想法是,当您尝试腌制实例方法时会导致错误,而您应该在这里尝试使用静态方法(使用实例方法似乎不正确,因为进程之间不共享对象)。
在声明 ParallelLoopTest 之前使用装饰器 @staticmethod 并删除 self 参数。

编辑: 另一种可能性是参数之一dd, final_col_dates, d1, a, d, final_col_dates_mod, iter, return_list 不能被取消腌制。显然,它来自panda.DataFrame
我看不出在这种情况下解封失败的任何原因,但我不太了解panda
一种解决方法是将数据转储到临时文件中。您可以查看此链接here 以有效地序列化panda.DataFrame。另一种解决方案是使用DataFrame.to_pickle 方法和panda.read_pickle 将其转储到文件/从文件中检索。

请注意,最好将joblib.Parallelmultiprocessing.Pool 进行比较,而不是与multiprocessing.Process 进行比较。

【讨论】:

  • 请看我最近的编辑。解决方案似乎不起作用。我是否以任何错误的方式实现它?
猜你喜欢
  • 1970-01-01
  • 2017-03-07
  • 2015-12-11
  • 2016-03-03
  • 2020-05-29
  • 1970-01-01
  • 1970-01-01
  • 2014-05-14
  • 2012-10-27
相关资源
最近更新 更多