You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让Python始终维持固定数量的线程运行?

如何让Python始终维持固定数量的线程运行?

我明白你的需求啦——你希望线程池里始终保持固定数量的线程在运行,而不是等第一批任务全跑完才启动新任务。咱们来看看你的代码问题出在哪,再一步步把它改好~

首先说下你现有代码的小问题:

  1. 代码里有个缩进错误,next_task = len(futures) + 1这一行的缩进不对,导致这部分逻辑根本跑不通;
  2. 更关键的是,你用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

这个代码的工作逻辑是:

  1. 先提交第一批max_workers个任务,用字典把future对象和任务编号绑定起来;
  2. 进入循环,用as_completed持续监控所有待完成的future;
  3. 每完成一个任务,就从字典里移除对应的future,然后检查是否还有任务需要运行,如果有,立刻提交新任务并加入字典;
  4. 不管任务是成功还是失败,都及时补上新任务,确保线程池里始终有max_workers个线程在跑。

如果你的需求是无限运行(比如持续从某个任务队列里取任务处理),只需要把total_tasks_to_run的判断逻辑改成从队列里取任务,比如:只要队列不为空,就提交新任务,这样就能一直维持固定的线程数量啦。

备注:内容来源于stack exchange,提问作者rkdnl rkdnl

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.14 09:03:05