Fra*_*ank 5 python multithreading
如何限制Python中并发线程的数量?
例如,我有一个包含许多文件的目录,我想处理所有这些文件,但是并行处理只有4个文件.
这是我到目前为止:
def process_file(fname):
# open file and do something
def process_file_thread(queue, fname):
queue.put(process_file(fname))
def process_all_files(d):
files=glob.glob(d + '/*')
q=Queue.Queue()
for fname in files:
t=threading.Thread(target=process_file_thread, args=(q, fname))
t.start()
q.join()
def main():
process_all_files('.')
# Do something after all files have been processed
Run Code Online (Sandbox Code Playgroud)
如何修改代码,以便一次只运行4个线程?
请注意,我想等待处理所有文件,然后继续处理已处理的文件.
例如,我有一个包含许多文件的目录,我想处理所有这些文件,但是并行处理只有4个文件.
这正是线程池的作用:您创建作业,并且池一次并行运行4个.您可以通过使用执行程序使事情变得更简单,您可以将其交给函数(或其他可调用函数),并将结果交给期货.你可以自己构建所有这些,但你没必要.*
stdlib的concurrent.futures模块是最简单的方法.(对于Python 3.1及更早版本,请参阅backport.)事实上,其中一个主要示例非常接近您想要做的事情.但是让我们根据您的具体用例进行调整:
def process_all_files(d):
files = glob.glob(d + '/*')
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
fs = [executor.submit(process_file, file) for file in files]
concurrent.futures.wait(fs)
Run Code Online (Sandbox Code Playgroud)
如果你想要process_file归还某些内容,那几乎就是这么简单:
def process_all_files(d):
files = glob.glob(d + '/*')
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor:
fs = [executor.submit(process_file, file) for file in files]
for f in concurrent.futures.as_completed(fs):
do_something(f.result())
Run Code Online (Sandbox Code Playgroud)
如果你也想处理异常......好吧,看看这个例子; 它只是呼叫的一个try/ except周围result().
*如果你想自己建造它们,那并不难.源代码multiprocessing.pool写得很好,评论也很复杂,并且大多数硬件与线程无关; 来源concurrent.futures更简单.