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

多线程及跨模块对象间共享PriorityQueue的实现问题

我来给你几个实用的解决方案,核心就是让不同模块的爬取、下载线程能安全共用同一个优先级队列,同时兼顾系统后续的扩展性:

方案1:依赖注入(最推荐,低耦合易扩展)

这是最符合软件工程设计原则的方式——顶层主模块初始化队列后,把队列作为参数“传递”给爬取、下载线程的实例。每个模块只需要专注于自己的业务逻辑,不用关心队列是在哪里创建的,后续新增任务线程时,直接传队列就行。

代码示例:

主模块 main.py

from queue import PriorityQueue
from crawler import CrawlerThread
from downloader import DownloaderThread
import threading

if __name__ == "__main__":
    # 顶层初始化优先级队列
    task_queue = PriorityQueue()

    # 把队列注入给爬取、下载线程实例
    crawler = CrawlerThread(task_queue)
    downloader = DownloaderThread(task_queue)

    # 启动线程
    crawler.start()
    downloader.start()

    # 按需处理线程退出(比如监听停止信号)
    # crawler.join()
    # downloader.join()

爬取模块 crawler.py

from queue import PriorityQueue
import threading
import time

class CrawlerThread(threading.Thread):
    def __init__(self, task_queue: PriorityQueue):
        super().__init__()
        self.task_queue = task_queue  # 保存队列引用

    def run(self):
        # 模拟持续爬取远程服务器的逻辑
        while True:
            # 替换为实际爬取逻辑:检测到新文件后生成任务信息
            priority = 1  # PriorityQueue中数字越小优先级越高
            file_task = {
                "remote_url": "http://your-server/new-file.zip",
                "local_path": "./downloads/new-file.zip"
            }
            self.task_queue.put((priority, file_task))
            print(f"新增下载任务: {file_task['remote_url']}")
            time.sleep(5)  # 模拟爬取间隔

下载模块 downloader.py

from queue import PriorityQueue
import threading
import time

class DownloaderThread(threading.Thread):
    def __init__(self, task_queue: PriorityQueue):
        super().__init__()
        self.task_queue = task_queue

    def run(self):
        while True:
            # 阻塞等待队列中的任务(线程安全)
            priority, file_task = self.task_queue.get()
            try:
                # 替换为实际下载逻辑:比如用requests/urllib下载文件
                print(f"开始下载[优先级{priority}]: {file_task['remote_url']}")
                time.sleep(3)  # 模拟下载耗时
                print(f"下载完成: {file_task['local_path']}")
            finally:
                # 标记任务完成,避免队列join时阻塞
                self.task_queue.task_done()

这个方案的优势:

  • 完全解耦,爬取/下载模块不依赖任何全局状态
  • 扩展性极强,新增任务线程时只需注入队列即可
  • 天然线程安全,PriorityQueue本身的put()/get()方法已经加了锁
方案2:全局模块变量(小型项目快速实现)

如果你的项目规模很小,想快速落地,可以把队列放在一个单独的共享模块里,其他模块直接导入使用。

代码示例:

共享队列模块 shared_queue.py

from queue import PriorityQueue

# 全局队列实例,模块导入时自动初始化
task_queue = PriorityQueue()

主模块 main.py

from shared_queue import task_queue
from crawler import CrawlerThread
from downloader import DownloaderThread
import threading

if __name__ == "__main__":
    crawler = CrawlerThread()
    downloader = DownloaderThread()

    crawler.start()
    downloader.start()

爬取/下载模块

直接导入from shared_queue import task_queue,然后调用task_queue.put()/task_queue.get()即可,逻辑和方案1一致。

注意:这种方式耦合度较高,后续如果需要修改队列的初始化逻辑(比如设置最大容量),会涉及多个模块的改动,适合小型项目快速验证。

方案3:单例模式(需控制队列初始化逻辑时用)

如果需要对队列的初始化做特殊配置(比如读取配置文件设置参数),可以把队列包装成单例类,确保全局只有一个实例。

代码示例:

单例队列模块 queue_singleton.py

from queue import PriorityQueue

class QueueSingleton:
    _instance = None

    def __new__(cls):
        if cls._instance is None:
            # 这里可以添加自定义初始化逻辑,比如读取配置
            cls._instance = PriorityQueue(maxsize=100)
        return cls._instance

# 导出全局单例实例
task_queue = QueueSingleton()

其他模块直接导入from queue_singleton import task_queue使用即可,和方案2的用法一致,但好处是可以统一控制队列的初始化逻辑。

额外注意事项
  • PriorityQueue本身是线程安全的,不需要你额外加锁,直接调用它的方法就行
  • 记得给线程添加退出逻辑,比如用threading.Event作为停止信号,避免线程无限循环跑下去

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 08:03:03