多线程及跨模块对象间共享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
相关产品推荐
相关产品推荐

