当multiprocessing.Queue达指定容量时,如何暂停线程任务添加?
解决任务队列堆积的最优方案
针对你遇到的生产者线程(读文件)速度远超消费者进程(爬站解析)导致的队列堆积问题,最可靠的解决方式是让生产者在队列达到阈值时自动阻塞,直到消费者处理完部分任务再恢复,下面是两种落地实现:
方案1:利用队列原生的maxsize阻塞(推荐)
multiprocessing.Queue支持设置maxsize参数,当队列里的任务数达到这个值时,调用put()会直接阻塞,直到队列有空闲位置——完全不用手动检查队列状态,代码最简洁:
from multiprocessing import Queue, Process import threading # 设定队列最大任务数N task_queue = Queue(maxsize=N) def read_subdomains(): with open('subdomains.txt', 'r') as f: for line in f: subdomain = line.strip() # 只处理目标平台的子域名 if subdomain.endswith(('.shopee.com', '.tokopedia.com')): # 队列满时自动停在这里,等消费者取走任务再继续 task_queue.put(subdomain) # 所有任务读完后,放一个终止标记告诉进程可以结束了 task_queue.put(None) def crawl_task(): while True: task = task_queue.get() if task is None: # 把终止标记放回队列,让其他进程也能收到 task_queue.put(None) break # 执行爬取和解析逻辑 crawl_and_parse(task) # 启动线程和进程 reader_thread = threading.Thread(target=read_subdomains) reader_thread.start() # 根据CPU核数或实际需求设置进程数 workers = [Process(target=crawl_task) for _ in range(4)] for worker in workers: worker.start() reader_thread.join() for worker in workers: worker.join()
方案2:用信号量做更灵活的阈值控制
如果不想严格等队列满才暂停,而是希望达到某个阈值(比如N的80%)就停,或者需要自定义阻塞逻辑,可以用线程信号量:
from multiprocessing import Queue, Process import threading task_queue = Queue() # 信号量初始值设为N,代表最多允许同时存在N个待处理任务 sem = threading.Semaphore(N) def read_subdomains(): with open('subdomains.txt', 'r') as f: for line in f: subdomain = line.strip() if subdomain.endswith(('.shopee.com', '.tokopedia.com')): # 获取信号量,余量为0时阻塞 sem.acquire() task_queue.put(subdomain) task_queue.put(None) def crawl_task(): while True: task = task_queue.get() if task is None: task_queue.put(None) break try: crawl_and_parse(task) finally: # 不管爬取成功失败,都释放信号量,避免死锁 sem.release() # 启动逻辑和方案1一致
关键注意点
- 方案1是官方推荐的方式,无需额外同步逻辑,出错概率最低,优先用这个。
- 方案2里的
sem.release()一定要放在finally块里,不然爬取过程抛出异常会导致信号量无法释放,生产者永远阻塞。 - 终止标记
None的传递要做好,确保所有工作进程都能收到并正常退出,避免出现僵尸进程。
内容的提问来源于stack exchange,提问作者MC874
相关产品推荐
相关产品推荐

