如何用multiprocessing实现生产者-消费者模式?含Manager()应用需求
当然可以用
Managers()实现这个需求! 你说得太对了,这完全是经典的生产者-消费者问题——生产者负责拉取文件存到本地,消费者负责把本地文件上传出去,核心就是要把本地目录的文件数控制在10个以内,防止磁盘过载。
为什么用Managers()?
multiprocessing.Managers()的优势在于能帮你在多进程之间轻松共享复杂的同步对象,比如带容量限制的队列、共享计数器或者锁。相比单纯用Queue,它的灵活性更高,尤其是当你需要跨进程共享状态(比如实时统计本地文件数量)的时候,用起来更顺手。
具体实现思路
我给你整理了一个贴合你需求的代码示例,核心逻辑是:
- 生产者进程:从你的
sha_list里逐个取文件标识,只有当本地“待上传队列”有空闲(也就是本地文件数小于10)时,才下载文件到本地。 - 消费者进程:持续从队列里取任务,上传对应文件,完成后删除本地文件,释放队列空间让生产者继续下载。
from multiprocessing import Process, Manager import requests import os import warnings # 只忽略requests的SSL警告,别全关哈 warnings.filterwarnings("ignore", category=requests.packages.urllib3.exceptions.InsecureRequestWarning) def producer(sha_list, file_queue, download_dir): os.makedirs(download_dir, exist_ok=True) for sha in sha_list: # 队列满时会自动阻塞,直到消费者取走一个任务(本地文件数降到9) file_queue.put(sha) # 替换成你实际的下载逻辑 file_path = os.path.join(download_dir, f"{sha}.bin") try: resp = requests.get(f"https://your-source-url/{sha}", verify=False) resp.raise_for_status() with open(file_path, "wb") as f: f.write(resp.content) print(f"已下载: {file_path}") except Exception as e: print(f"下载失败 {sha}: {str(e)}") # 下载失败要把sha从队列里移除,避免消费者空等 file_queue.get() file_queue.task_done() def consumer(file_queue, download_dir, upload_url): while True: sha = file_queue.get() # 接收到结束信号就退出 if sha is None: file_queue.task_done() break file_path = os.path.join(download_dir, f"{sha}.bin") if not os.path.exists(file_path): print(f"文件不存在,跳过: {file_path}") file_queue.task_done() continue # 替换成你实际的上传逻辑 try: with open(file_path, "rb") as f: resp = requests.post(upload_url, files={"file": f}, verify=False) resp.raise_for_status() os.remove(file_path) print(f"已上传并删除: {file_path}") except Exception as e: print(f"上传失败 {sha}: {str(e)}") # 可以根据需求添加重试逻辑 finally: file_queue.task_done() if __name__ == "__main__": sha_list = ["sha-001", "sha-002", "sha-003"] # 你的完整sha列表 download_dir = "./temp_downloads" upload_target_url = "https://your-target-url.com/upload" with Manager() as manager: # 队列maxsize设为10,刚好对应允许的最大本地文件数 file_queue = manager.Queue(maxsize=10) # 启动进程,你可以根据需求调整消费者数量(比如开2个消费者加快上传) producer_proc = Process(target=producer, args=(sha_list, file_queue, download_dir)) consumer_proc = Process(target=consumer, args=(file_queue, download_dir, upload_target_url)) producer_proc.start() consumer_proc.start() # 等生产者完成所有下载任务 producer_proc.join() # 给消费者发结束信号 file_queue.put(None) consumer_proc.join() print("所有任务完成!")
关键细节说明
- 队列容量控制:
manager.Queue(maxsize=10)是核心,队列满的时候生产者会自动阻塞,不用你手动统计本地文件数——队列的长度就对应着本地待上传的文件数。 - 异常处理:我加了下载/上传失败的捕获逻辑,避免单个任务失败导致整个流程卡住。下载失败时要记得把任务从队列里移除,上传失败可以根据需求加重试。
- 进程同步:这里用队列的
task_done()和join()来保证所有任务都被处理,不需要额外加锁(队列本身是线程/进程安全的),除非你需要更复杂的状态共享。
替代方案?
其实不用Managers(),单纯用multiprocessing.Queue也能实现类似的容量控制,但Managers()更适合需要共享多个同步对象的场景(比如同时共享计数器和锁),而且跨进程共享的稳定性更好。
内容的提问来源于stack exchange,提问作者Ishan Bhatt
相关产品推荐
相关产品推荐

