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

如何等待线程执行完毕后启动新线程,实现固定并发量的多线程调度?

问题根因

你的代码丢失后续任务的核心原因是:for i in range(total_runs) 循环每次迭代会固定推进任务编号i,当活跃线程数超过阈值时,本次循环仅清理已结束的线程、不会启动当前i对应的任务,直接进入下一轮循环,相当于直接丢弃了当前i的任务。

方案1:最小改动修复原有逻辑

只要在线程满时,阻塞等待直到有空闲线程,再启动当前任务即可,修改后代码如下:

import random
import time
import threading

def random_wait(thread_id):
    rand_number = random.randint(3, 9)
    print(f'Thread {thread_id} waiting for {rand_number} seconds')
    time.sleep(rand_number)
    print(f'Thread {thread_id} completed execution')

if __name__=='__main__':
    total_runs = 6
    thread_limit = 3
    running_threads = list()
    i = 0
    while i < total_runs:
        # 先清理所有已结束的线程
        running_threads = [t for t in running_threads if t.is_alive()]
        # 有空闲容量就启动新线程
        if len(running_threads) < thread_limit:
            thread = threading.Thread(target=random_wait, args=(i,))
            running_threads.append(thread)
            print(f'Starting thread : {i}')
            thread.start()
            i += 1
        else:
            # 没有空位就稍等100ms再检查,避免空耗CPU
            time.sleep(0.1)
    # 等待所有线程执行完成
    for t in running_threads:
        t.join()
方案2:更优方案1:使用标准库线程池(最推荐)

Python 标准库concurrent.futures.ThreadPoolExecutor已经封装了线程数控制、空闲线程自动分配任务的能力,不需要手动管理线程生命周期,代码更简洁可靠:

import random
import time
from concurrent.futures import ThreadPoolExecutor

def process_folder(folder_id):
    # 这里替换为实际的文件夹处理逻辑
    rand_number = random.randint(3, 9)
    print(f'Processing folder {folder_id}, will take {rand_number}s')
    time.sleep(rand_number)
    print(f'Folder {folder_id} processed done')

if __name__=='__main__':
    total_folders = 100
    max_threads = 10
    # 固定最大线程数为10
    with ThreadPoolExecutor(max_workers=max_threads) as executor:
        # 提交所有100个文件夹的处理任务,线程池会自动分配空闲线程执行
        for folder_id in range(total_folders):
            executor.submit(process_folder, folder_id)

这个实现会自动保持最多10个线程同时运行,只要有线程执行完,立刻分配下一个待处理的文件夹任务,完全符合你的需求。

方案3:更优方案2:固定线程+任务队列(无额外线程创建开销)

如果你要严格固定启动10个线程、不动态创建销毁线程,可以用队列存储待处理任务,10个线程循环从队列取任务执行:

import random
import time
import threading
from queue import Queue

def worker(queue):
    while not queue.empty():
        # 从队列取待处理的文件夹ID
        folder_id = queue.get()
        try:
            # 这里替换为实际的文件夹处理逻辑
            rand_number = random.randint(3, 9)
            print(f'Thread {threading.current_thread().name} processing folder {folder_id}, will take {rand_number}s')
            time.sleep(rand_number)
            print(f'Thread {threading.current_thread().name} finished folder {folder_id}')
        finally:
            # 标记任务完成
            queue.task_done()

if __name__=='__main__':
    total_folders = 100
    max_threads = 10
    # 创建任务队列,存入所有待处理的文件夹ID
    task_queue = Queue()
    for folder_id in range(total_folders):
        task_queue.put(folder_id)
    # 启动固定10个工作线程
    for i in range(max_threads):
        t = threading.Thread(target=worker, args=(task_queue,), name=f"Worker-{i}")
        t.start()
    # 等待所有任务处理完成
    task_queue.join()

这个方案线程只会创建一次,运行过程中不会有新线程创建销毁的开销,更适合长时间运行的计算密集/GPU任务场景。

内容的提问来源于stack exchange,提问作者Aditya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 11:48:03