如何让Python始终维持固定数量的线程运行?
如何让Python始终维持固定数量的线程运行?
我明白你的需求啦——你希望线程池里始终保持固定数量的线程在运行,而不是等第一批任务全跑完才启动新任务。咱们来看看你的代码问题出在哪,再一步步把它改好~
首先说下你现有代码的小问题:
- 代码里有个缩进错误,
next_task = len(futures) + 1这一行的缩进不对,导致这部分逻辑根本跑不通; - 更关键的是,你用
concurrent.futures.as_completed(futures)的时候,这个函数只会遍历你一开始传入的那10个future对象,后续添加到futures列表里的新任务不会被这个循环自动监控到,所以等第一批10个任务里的几个完成后,新提交的任务完成时,不会触发新的任务提交,线程池最后还是会空下来。
那怎么改才能让线程池始终维持固定数量的线程呢?核心思路是:持续监控所有待完成的任务(包括新提交的),每完成一个就立刻补上一个新任务。
给你一个修正后的示例代码,我还调整了任务的睡眠时间,方便你更直观看到效果:
import concurrent.futures import time def example_task(n): print(f"Task {n} started.") # 把睡眠时间改成随机一点的,避免任务同时完成,方便观察线程维持效果 time.sleep(n % 3 + 1) print(f"Task {n} completed.") return n max_workers = 10 # 这里可以设置你要运行的总任务数,或者改成无限循环的逻辑 total_tasks_to_run = 20 completed_tasks = 0 with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor: # 用字典存future和对应的任务编号,方便后续追踪 futures = {executor.submit(example_task, i+1): i+1 for i in range(max_workers)} # 持续处理完成的任务,直到所有任务都跑完 while futures: # 等待任意一个任务完成 for future in concurrent.futures.as_completed(futures): task_num = futures.pop(future) completed_tasks += 1 try: result = future.result() print(f"Result of task {task_num}: {result}") # 如果还有任务要跑,立刻提交新任务补位 if completed_tasks < total_tasks_to_run: next_task_num = completed_tasks + 1 new_future = executor.submit(example_task, next_task_num) futures[new_future] = next_task_num except Exception as e: print(f"Task {task_num} failed with error: {e}") # 就算任务出错了,也要补上新任务维持线程数量 if completed_tasks < total_tasks_to_run: next_task_num = completed_tasks + 1 new_future = executor.submit(example_task, next_task_num) futures[new_future] = next_task_num
这个代码的工作逻辑是:
- 先提交第一批
max_workers个任务,用字典把future对象和任务编号绑定起来; - 进入循环,用
as_completed持续监控所有待完成的future; - 每完成一个任务,就从字典里移除对应的future,然后检查是否还有任务需要运行,如果有,立刻提交新任务并加入字典;
- 不管任务是成功还是失败,都及时补上新任务,确保线程池里始终有
max_workers个线程在跑。
如果你的需求是无限运行(比如持续从某个任务队列里取任务处理),只需要把total_tasks_to_run的判断逻辑改成从队列里取任务,比如:只要队列不为空,就提交新任务,这样就能一直维持固定的线程数量啦。
备注:内容来源于stack exchange,提问作者rkdnl rkdnl
相关产品推荐
相关产品推荐

