Python多线程队列应用:未知任务数下限制10并发线程的实现咨询
解决方案:用任务队列+固定工作线程池控制并发
针对你的需求,我们可以用Python标准库中的queue.Queue实现生产者-消费者模式:让10个工作线程持续从队列里取任务执行,读取文件的主线程负责把收集好的big_list作为任务放进队列,这样就能天然控制同时运行的线程数不超过10个。
下面是修改后的完整代码:
import queue import threading # 线程安全的全局任务队列 task_queue = queue.Queue() # 最大并发工作线程数 MAX_WORKERS = 10 def worker(): """持续运行的工作线程:从队列取任务执行,直到收到停止信号""" while True: task_data = task_queue.get() # 收到None代表需要停止工作 if task_data is None: task_queue.task_done() break # 执行实际的记录处理逻辑 processList(task_data) task_queue.task_done() def process_data(line, big_list): """修改后的收集逻辑:完成一条记录就放入任务队列""" # 注意strip()避免换行符干扰判断 if line.strip() == '***RECORD END***': # 放入队列前拷贝列表,避免后续修改影响已入队的任务 task_queue.put(big_list.copy()) big_list.clear() # 重置列表准备收集下一条记录 return big_list.append(line) if __name__ == "__main__": # 启动10个工作线程 for _ in range(MAX_WORKERS): t = threading.Thread(target=worker) t.daemon = True # 可选:主线程结束时自动终止工作线程 t.start() big_list = [] # 读取大文件并收集数据 with open('verybigfile.txt', 'r') as f: for line in f: process_data(line, big_list) # 处理文件末尾可能未完成的最后一条记录 if big_list: task_queue.put(big_list) # 等待队列中所有任务处理完成 task_queue.join() # 给每个工作线程发送停止信号 for _ in range(MAX_WORKERS): task_queue.put(None)
关键逻辑说明:
- 任务队列
queue.Queue:它本身是线程安全的,自动处理任务分发和同步,确保同一时刻最多有10个线程在执行任务(因为只有10个工作线程在取任务)。 - 工作线程
worker:启动后会一直循环从队列取任务,直到收到None才退出,完美实现了"持续运行"的需求。 - 数据收集逻辑:每次收集完一条完整记录,就把
big_list的副本放入队列(用copy()避免后续修改影响已入队的任务),然后重置列表。 - 收尾处理:文件读取完成后,先调用
task_queue.join()等待所有任务处理完毕,再给队列放入和工作线程数相同的None,让每个线程都能收到停止信号并安全退出。
注意点:
- 如果
processList会修改全局变量,记得用threading.Lock()加锁保证线程安全。 - 守护线程的设置:如果你的程序还有其他非守护线程需要运行,可以去掉
daemon=True;如果希望主线程结束时工作线程自动退出,保留该设置即可。
内容的提问来源于stack exchange,提问作者imtephi
相关产品推荐
相关产品推荐

