如何用Python Queue分批处理大型线程列表?每次限10个线程
用Python Queue控制线程并发数(每次仅运行10个)
嘿,这个场景我太熟悉了!当你有200+线程要处理时,直接全部启动很容易把系统资源拉满,用queue.Queue来做并发控制是非常靠谱的方案。我给你两种思路,根据你的实际情况选就行:
推荐方案:任务队列+固定消费者线程(更高效)
比起预先创建200多个线程,我更建议你把任务逻辑放进队列,然后启动10个固定的工作线程来消费任务——这样线程复用率更高,系统资源占用更稳定。
示例代码
import queue import threading import time # 定义你的任务逻辑(替换成你实际要执行的代码) def process_task(task_id): print(f"启动任务 {task_id}") # 模拟耗时操作,比如接口请求、文件处理等 time.sleep(2) print(f"完成任务 {task_id}") def main(): # 1. 创建线程安全的任务队列 task_queue = queue.Queue() # 2. 往队列中添加所有任务(这里模拟200个任务) total_tasks = 200 for task_id in range(total_tasks): task_queue.put(task_id) # 3. 定义工作线程的逻辑:不断从队列取任务执行 def worker_thread(): while not task_queue.empty(): try: # 非阻塞获取任务,避免队列空时一直等待 current_task = task_queue.get(block=False) process_task(current_task) task_queue.task_done() # 标记任务完成(配合join使用) except queue.Empty: break # 4. 启动10个工作线程 workers = [] for _ in range(10): thread = threading.Thread(target=worker_thread) thread.start() workers.append(thread) # 5. 等待所有工作线程完成 for worker in workers: worker.join() # 可选:等待队列中所有任务标记完成 task_queue.join() print("所有任务处理完毕!") if __name__ == "__main__": main()
为什么推荐这个方案?
- 不用一次性创建200多个线程,复用10个线程就能处理所有任务,内存开销小很多
queue.Queue是线程安全的,不用自己加锁处理并发问题- 任务的添加和消费逻辑解耦,后续调整并发数(比如改成15个)只需要改一行代码
备选方案:用Queue控制已有线程列表的并发
如果你已经预先创建了200+线程对象(threadlist),可以用一个容量为10的Queue作为许可池,每个线程必须拿到许可才能执行核心逻辑,执行完后归还许可,以此控制同时运行的线程数。
示例代码
import queue import threading import time # 假设这是你每个线程要执行的核心逻辑 def thread_core_logic(thread_id, permission_queue): # 获取运行许可:队列满时会自动等待,直到有许可释放 permission_queue.get() try: print(f"线程 {thread_id} 开始运行") time.sleep(2) # 替换成你的实际业务代码 print(f"线程 {thread_id} 运行结束") finally: # 必须在finally里归还许可,避免线程异常导致许可丢失 permission_queue.put(1) def main(): # 创建容量为10的许可队列,预先放入10个"许可" permission_queue = queue.Queue(maxsize=10) for _ in range(10): permission_queue.put(1) # 构建你的线程列表(模拟200个线程) threadlist = [] total_threads = 200 for thread_id in range(total_threads): thread = threading.Thread( target=thread_core_logic, args=(thread_id, permission_queue) ) threadlist.append(thread) # 启动所有线程 for thread in threadlist: thread.start() # 等待所有线程执行完成 for thread in threadlist: thread.join() print("所有线程运行完毕!") if __name__ == "__main__": main()
注意事项
- 一定要在
finally块中归还许可,否则如果某个线程抛出异常,许可会丢失,导致后续线程一直等待 - 这种方案适合你已经有现成线程对象的场景,但线程复用率不如第一种方案高
额外小建议
- 如果你的任务是IO密集型(比如爬虫、接口调用),10个并发数是合理的,甚至可以适当调高;如果是CPU密集型任务,建议并发数不要超过你的CPU核心数,避免频繁上下文切换拖慢速度
- 如果你需要更强大的并发控制,可以考虑用
concurrent.futures.ThreadPoolExecutor(本质也是基于线程池,用法更简洁),不过既然你指定要用Queue,上面的方案完全能满足需求
内容的提问来源于stack exchange,提问作者Sid
相关产品推荐
相关产品推荐

