【问题标题】:Parallel processing a function that's in a separate module并行处理单独模块中的函数
【发布时间】:2014-01-13 12:39:40
【问题描述】:

我有一个“令人尴尬的并行”任务:我正在尝试以 CPU 繁重的方式解析大量日志文件。我不关心它们完成的顺序,进程不需要共享任何资源或线程。

我在 Windows 机器上。

我的设置是这样的:

main.py

import parse_file
import multiprocessing

...

files_list = ['c:\file1.log','c:\file2.log']

if __name__ == '__main__':
    pool = multiprocessing.Pool(None)

    for this_file in files_list:
        r = pool.apply_async(parse_file.parse, (this_file, parser_config))
        results = r.get()

...

#Code to do stuff with the results

parse_file 基本上是一个完全独立的模块,不访问任何共享资源 - 结果以列表形式返回。

当我在没有多处理的情况下运行它时,这一切都运行得很好,但是当我启用它时,会发生一堵巨大的错误墙,表明源模块(其中的那个)是正在运行的模块并行运行。 (该错误是仅在源脚本(不是 parse_file 模块)中的数据库锁定错误,并且在多处理之前的某个点!)

我并没有假装理解多处理模块,而是从 other 示例 here 开始工作,但没有一个包含任何表明这是正常的或为什么会发生的内容。

我做错了什么?如何多处理此任务? 谢谢!


使用此方法可轻松复制: 测试.py

import multiprocessing
import test_victim

files_list = ['c:\file1.log','c:\file2.log']

print("Hello World")

if __name__ == '__main__':
    pool = multiprocessing.Pool(None)
    results = []
    for this_file in files_list:
        r = pool.map_async(test_victim.calculate, range(10), callback=results.append)
        results = r.get()

    print(results)

test_victim.py:

def calculate(value):
    return value * 10

运行 test.py 时的输出应该是:

Hello World
[0, 10, 20, 30, 40, 50, 60, 70, 80, 90]

但实际上是这样的:

Hello World
[0, 10, 20, 30, 40, 50, 60, 70, 80, 90]
Hello World
Hello World

(额外的“Hello World”的实际数量)每次我运行它时都会在 1 到 4 之间变化 = 应该没有)

【问题讨论】:

  • 确保for-loop 缩进,以便在if __name__ ... 语句内。否则代码将在 Windows 上导入炸弹。
  • @unutbu - 谢谢。根据您的建议,我刚刚这样做了,但恐怕这完全没有区别。 :-((更新示例以反映这一点))
  • 请发布堆栈跟踪,至少前几行和最后几行。
  • 在堆栈跟踪中,以File 开头的最后一行是指脚本的路径(main.py)是什么?后面的线是什么?
  • 我不知道问题出在哪里,但在我看来,我们需要了解您使用 sqlite3 的结构。重现错误的可运行示例将非常棒。

标签: python python-3.x multiprocessing


【解决方案1】:

在 Windows 上,当 Python 执行时

pool = multiprocessing.Pool(None) 

产生了新的 Python 进程。因为 Windows 没有os.fork,所以这些新的 Python 进程重新导入调用模块。因此,任何不在里面的东西

if __name__ == '__main__': 

为每个生成的进程执行一次。这就是为什么您会看到多个Hello Worlds

请务必阅读文档中的 "Safe importing of main module" 警告。


所以要修复,将所有只需要运行一次的代码放入

if __name__ == '__main__': 

声明。


例如,您的可运行示例将通过放置来修复

print("Hello World")

if __name__ == '__main__' 语句中:

import multiprocessing
import test_victim

files_list = ['c:\file1.log','c:\file2.log']

def main():
    print("Hello World")
    pool = multiprocessing.Pool(None)
    results = []
    for this_file in files_list:
        r = pool.map_async(test_victim.calculate, range(10), callback=results.append)
        results = r.get()

    print(results)

if __name__ == '__main__':
    main()

产量

Hello World
[0, 10, 20, 30, 40, 50, 60, 70, 80, 90]

尤其是在 Windows 上,使用 multiprocessing 的脚本必须既可运行(作为脚本)又可导入。使脚本可导入的一种简单方法是按如上所示对其进行结构化。将 脚本 应该执行的所有内容放在一个名为 main 的函数中,然后使用

if __name__ == '__main__':
    main()

在脚本的末尾。 main 之前的东西应该只是 import 语句和全局常量的定义。

【讨论】:

  • 谢谢。在发布这个问题之前,我实际上尝试阅读该文档,但它假设我根本没有的知识水平(即,它没有意义)。 main.py 太复杂,无法添加 name 东西。所以我创建了一个仅包含多处理内容的单独文件,并从 main.py 调用它,但我仍然无法让它工作(如果没有运行 name 中的任何内容 永远 - 如果它被遗漏,我会得到各种随机性)。你能提供一个test.py的例子吗?
  • 谢谢!最后,我只是将if __name__ == '__main__': 作为我的 main.py 脚本的第一行。现在一切都是丑陋的缩进,但是很好。再次感谢。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2015-03-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多