如何等待线程执行完毕后启动新线程,实现固定并发量的多线程调度?
问题根因
你的代码丢失后续任务的核心原因是: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
相关产品推荐
相关产品推荐

