【问题标题】:multiprocessing.Pool.map does not work in parallelmultiprocessing.Pool.map 不能并行工作
【发布时间】:2018-05-06 17:43:10
【问题描述】:

我想要达到的目标:

并行化一个每次调用产生多个线程的函数,如下所示:

 - PROCESS01 -> 16 Threads
 - PROCESS02 -> 16 Threads
 - ...
 - PROCESSn -> 16 Threads

代码:

with multiprocessing.Pool(4) as process_pool:
    results = process_pool.map(do_stuff, [drain_queue()])

drain_queue() 返回项目列表和

do_stuff(item_list):
    print('> PID: ' + str(os.getpid()))
    with concurrent.futures.ThreadPoolExecutor(max_workers=16) as executor:
        result_dict = {executor.submit(thread_function, item): item for item in item_list}
        for future in concurrent.futures.as_completed(result_dict):
            pass

thread_function() 处理传递给它的每个项目。

但是,当执行代码时,输​​出如下:

> PID: 1000
(WAITS UNTIL THE PROCESS FINISHES, THEN START NEXT)
> PID: 2000
(WAITS UNTIL THE PROCESS FINISHES, THEN START NEXT)
> PID: 3000
(WAITS UNTIL THE PROCESS FINISHES, THEN START NEXT)
> PID: 3000
(WAITS UNTIL THE PROCESS FINISHES, THEN START NEXT)

Here is a screenshot of Task Manager

我在这里缺少什么?我无法弄清楚为什么不能按预期工作。 谢谢!

【问题讨论】:

  • 线程和进程是完全不同的东西。我不确定你为什么要混合它们。 Global Interpreter Lock (GIL) 确保一次只有一个线程可以执行字节码
  • 那么串行执行是由GIL引起的,它迫使每个进程等待线程执行?
  • 其实这是个好问题。启动多处理后,我不确定跨进程的线程操作是否应该相互干扰,如果它们是单独进程的一部分。我推迟回答,我无法明确回答。

标签: python python-3.x python-multiprocessing


【解决方案1】:

我发现了问题。 map() 的第二个参数应该是一个可迭代的,在我的例子中是一个包含一个 单个对象的列表。

什么错了吗?这个:[drain_queue()],它会生成一个包含单个对象的列表。

在这种情况下,代码

with multiprocessing.Pool(4) as process_pool:
    results = process_pool.map(do_stuff, [drain_queue()])

强制multiprocessing.Pool.map 将单个对象“分发”到单个进程,即使它创建n 数量的进程,工作仍将由一个进程完成。幸好与 GIL 限制无关。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-18
    • 2017-05-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多