You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

当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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 04:15:37