如何修改Python代码以限制最大并发线程数为5?
实现最大5个活跃线程的两种方案
你的需求是限制同时运行的线程数为5,线程完成后自动启动下一个任务,以下两种方法可以实现:
方案一:使用ThreadPoolExecutor(推荐,代码更简洁)
Python3.2+自带的concurrent.futures.ThreadPoolExecutor可以直接管理线程池,指定max_workers=5就能限制最大活跃线程数,它会自动处理线程的创建、复用和任务调度。
修改后的main函数代码:
import concurrent.futures def main(): get_emails() log("Starting tasks", "yellow") # 创建线程池,最大5个工作线程 with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor: # 给每个用户任务提交到线程池 executor.map(targetfunc, users.values()) log("All tasks completed!", "success")
说明:
executor.map会自动遍历users.values(),为每个元素启动一个任务,线程池会控制同时运行的线程数不超过5with语句会自动关闭线程池,等待所有任务完成后再执行后续代码
方案二:使用Semaphore手动控制线程数
如果你想手动控制线程的创建和并发数,可以用threading.Semaphore,它相当于一个"许可证",每次启动线程前获取许可证,线程结束后释放许可证,确保同时最多只有5个线程持有许可证。
步骤1:修改targetfunc,添加释放操作
import threading # 全局定义信号量,最多允许5个线程同时运行 semaphore = threading.Semaphore(5) def targetfunc(data): global users_done global successCount try: # 邮件提交逻辑保持不变 email = data["EMAIL"] # ... 你的原有代码 ... finally: # 线程完成后释放许可证,让新线程可以启动 semaphore.release()
步骤2:修改main函数,启动线程前获取许可证
def main(): get_emails() log("Starting tasks", "yellow") threads = [] for user in users: # 获取许可证,没有可用时会阻塞,直到有线程释放 semaphore.acquire() thread = threading.Thread(target=targetfunc, args=(users[user],)) thread.start() threads.append(thread) # 等待所有线程完成 for thread in threads: thread.join() log("All tasks completed!", "success")
注意:
- 把
semaphore.release()放在finally块里,确保即使任务执行出错,许可证也能被释放,避免死锁 - 如果
users_done或successCount是多线程修改的全局变量,建议用threading.Lock来保护,避免数据竞争:count_lock = threading.Lock() # 在修改变量时使用锁: with count_lock: successCount += 1
内容的提问来源于stack exchange,提问作者mxxim
相关产品推荐
相关产品推荐

