【问题标题】:How to limit the nr of threads in a thread pool for infinite iterable?如何限制线程池中线程的 nr 以实现无限迭代?
【发布时间】:2019-09-12 20:09:28
【问题描述】:

假设我们想对所有自然数进行无限迭代,并为每个数字打开一个线程,直到一个限制。由于有无限的自然数,这个列表就像一个生成器,那么我们如何在生成数字的同时限制打开线程的数量呢?

我正在尝试类似的东西:

from multiprocessing import Pool

def func(i):
    """ Do something with i"""

p = Pool(10)
i = 0
while True:
    p.imap_unordered(func,i)
    i += 1

我尝试了许多其他不同的方法(线程池、信号量等),但它们似乎都忽略了最大线程数。我想让一些东西打开 MAX_THREADS 个线程,每次线程完成时,它都会再次迭代并开始一个新线程。我怎样才能做到这一点?

【问题讨论】:

  • 您已经在池中设置了最大进程数。您根本没有使用线程。是什么让您认为您的设置被忽略并产生了任意数量的线程?

标签: python multithreading


【解决方案1】:

您想使用ThreadPoolExecutor。构造函数的第一个参数是max_workers,它决定了它可以使用多少个线程。

所以你最终会得到这样的结果:

with ThreadPoolExecutor(max_workers=n) as pool:
    future = pool.submit(function, *args)
future.result() # outputs result

返回值是Future,所以你可以在future完成后访问结果。当with 块结束时,ThreadPoolExecutor 等待其所有工作人员完成,因此您可以在块结束后安全地访问结果。


下面是ThreadPoolExecutor 的一个简单示例:

from concurrent.futures import ThreadPoolExecutor
from time import sleep

def fn(delay, result):
    sleep(x)
    print(result)

with ThreadPoolExecutor(max_workers=2) as pool:
    pool.submit(fn, 3, "last")
    pool.submit(fn, 1, "first")
    pool.submit(fn, 1, "middle")
# This should output "first", then "middle", then "last"

【讨论】:

  • 我已经试过了。我在任何地方都看不到 func 的输出。看起来它正在无限地迭代数字
  • @GuilhermeLima Executor.submit 返回一个Future,所以如果你说x = pool.submit(...),你可以在线程完成其工作后以x.result() 访问结果。如果您希望结果自动填充数据结构,您可以将它们保存在 Queue 中,而不是从函数中返回它们。
  • 我在工作函数中有一些打印语句,但它们没有被打印出来。怎么回事?
  • @GuilhermeLima 如果没有更多详细信息,我无法确定,但我添加了一个对我有用的示例来回答我的问题。它对你有用吗?
  • 我认为 OP 的主要观点是能够无限期地处理这些数字,并在一个完成后立即安排一个新工人
【解决方案2】:

使用 asyncio 的轻量级解决方案

import asyncio
import itertools as it
from random import randint


MAX_WORKERS = 5

async def worker(num):
    s = f"worker for {num}"
    print(s, "start")
    await asyncio.sleep(randint(1,2))
    print(s, "end")
    return num*100

async def main():
    seq = it.count()  # counts natural numbers endlessly
    pending = []  # the tasks which are pending
    while True:
        groups = zip(*[seq]*(MAX_WORKERS-len(pending)))  # fetching only as many inputs as the completed ones
        works = [
            worker(i) for i in next(groups)
        ]  # create as many workers as the ones which have completed

        # now await the completion of the first among the pending tasks and the new ones
        done, pending = await asyncio.wait(
            it.chain(pending, works),
            return_when=asyncio.FIRST_COMPLETED
        )
        print("Done", len(done), "Pending", len(pending))
        for t in done:  # print the result of each task which has completed
            print(t.result())


asyncio.run(main())

产生,例如:

worker for 1 start
worker for 4 start
worker for 0 start
worker for 2 start
worker for 3 start
worker for 1 end
worker for 4 end
worker for 0 end
worker for 2 end
Done 4 Pending 1
100
400
0
200
worker for 5 start
worker for 8 start
worker for 6 start
worker for 7 start
worker for 3 end
worker for 8 end
worker for 6 end
Done 3 Pending 2
300
800
600
worker for 11 start
worker for 10 start
worker for 9 start
worker for 5 end
worker for 7 end
worker for 11 end
worker for 9 end
Done 4 Pending 1
700
500
1100
900
worker for 15 start
worker for 12 start
worker for 13 start
worker for 14 start
worker for 10 end
Done 1 Pending 4
1000
worker for 16 start
worker for 15 end
worker for 13 end
worker for 14 end
Done 3 Pending 2
1400
1300
1500
worker for 18 start
worker for 17 start
worker for 19 start
worker for 16 end
worker for 12 end
Done 2 Pending 3
1600
1200
worker for 20 start
worker for 21 start
worker for 19 end
Done 1 Pending 4
1900
worker for 22 start
worker for 18 end
worker for 17 end
worker for 20 end
Done 3 Pending 2
1800
2000
1700
worker for 23 start
worker for 24 start
worker for 25 start
worker for 21 end
worker for 22 end
Done 2 Pending 3
2100
2200
worker for 27 start
worker for 26 start
worker for 23 end
worker for 24 end
worker for 25 end
worker for 26 end
Done 4 Pending 1
2500
2300
2600
2400
worker for 30 start
worker for 31 start
worker for 28 start
worker for 29 start
worker for 27 end
Done 1 Pending 4
2700
worker for 32 start
worker for 30 end
worker for 31 end
worker for 28 end
worker for 29 end
Done 4 Pending 1
2900
3000
2800
3100
worker for 33 start
worker for 34 start
worker for 35 start
worker for 36 start
worker for 32 end
worker for 35 end
worker for 36 end
Done 3 Pending 2
3600
3500
3200
worker for 39 start
worker for 38 start
worker for 37 start
worker for 33 end
worker for 34 end
Done 2 Pending 3
3400
3300
worker for 40 start
worker for 41 start
...

【讨论】:

    猜你喜欢
    • 2017-06-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-02-03
    相关资源
    最近更新 更多