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

Python多线程队列设置遇阻塞:如何实现有限队列并喂入任务?

问题根源

你现在的程序卡住,是因为先把所有URL全塞进队列,再启动线程消费——当队列设了maxsize=10,塞到第11个URL时队列满了,queue.put()会阻塞,而此时还没有任何线程在消费队列元素,程序就彻底卡在这里了。

解决思路

  1. 把「加载URL到队列」的操作放到独立线程里,让它和消费线程同时运行:队列满时put会暂时阻塞,等消费线程取走元素后再继续塞,避免一次性加载大量URL占用内存。
  2. 用「死亡药丸」机制让线程正确退出:所有URL加载完成后,给队列放和线程数相同的None(终止标记),每个线程拿到这个标记就停止工作。

修改后的完整代码

from pathlib import Path
from threading import Thread
from queue import Queue


class UrlConverter:
    def _load_urls(self, filename: str, queue: Queue):
        urls_file_path = str(Path(__file__).parent / Path(filename))
        with open(urls_file_path, 'r', encoding="utf-8") as txt_file:
            for line in txt_file:
                line = line.strip()
                if line:  # 跳过空行
                    queue.put(line)

    def start_loading(self, filename: str, queue: Queue) -> Thread:
        # 启动单独线程加载URL到队列
        load_thread = Thread(target=self._load_urls, args=(filename, queue))
        load_thread.start()
        return load_thread


class Fetcher:
    def worker(self, queue: Queue):
        while True:
            url = queue.get()
            # 拿到死亡药丸,结束工作
            if url is None:
                queue.task_done()
                break
            try:
                print(f"{url}: OK\n")
            except Exception as e:
                print(f"Error {url}: {str(e)}")
            finally:
                queue.task_done()  # 标记当前任务处理完成


class FetcherThreads:
    def __init__(self, thread_count: int):
        self.thread_count = thread_count
        self.threads = []

    def start_workers(self, queue: Queue):
        # 创建指定数量的工作线程
        for _ in range(self.thread_count):
            fetcher = Fetcher()
            worker_thread = Thread(target=fetcher.worker, args=(queue,), daemon=True)
            worker_thread.start()
            self.threads.append(worker_thread)

    def wait_for_finish(self, queue: Queue):
        # 等待所有URL处理完成
        queue.join()
        # 给每个线程发死亡药丸,通知退出
        for _ in range(self.thread_count):
            queue.put(None)
        # 等待所有工作线程结束
        for thread in self.threads:
            thread.join()


def main():
    MAX_QUEUE_SIZE = 10
    WORKER_THREADS = 10

    # 初始化带容量限制的队列
    urls_queue = Queue(maxsize=MAX_QUEUE_SIZE)

    # 启动URL加载线程
    url_converter = UrlConverter()
    load_thread = url_converter.start_loading('urls.txt', urls_queue)

    # 启动工作线程池
    thread_pool = FetcherThreads(WORKER_THREADS)
    thread_pool.start_workers(urls_queue)

    # 等待URL加载完成
    load_thread.join()

    # 等待所有任务处理完毕,终止线程
    thread_pool.wait_for_finish(urls_queue)


if __name__ == "__main__":
    main()

关键细节说明

  • 加载线程和消费线程并行运行,队列满时自动阻塞加载逻辑,直到有线程取走元素,完美控制内存占用。
  • 工作线程持续从队列取任务,直到拿到None才退出,不用反复创建销毁线程,效率更高。
  • 用queue.task_done()和queue.join()跟踪任务完成状态,确保所有URL都被处理。
  • 工作线程设为daemon=True,防止程序意外挂起,同时通过显式join保证线程正常退出。

内容的提问来源于stack exchange,提问作者zaelcovsky

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 08:22:45