Python多进程实现多生产者多消费者模型出现死锁该如何解决?
Python multiprocessing实现多生产者多消费者模型的问题与解决
背景与初始实现
我一直在尝试用Python的multiprocessing库实现多生产者多消费者模型,生产者负责从网页爬取数据,消费者负责处理数据。
一开始我仅实现了具备对应功能的生产者、消费者两个函数,使用Queue实现二者间通信,但一直不知道如何处理任务完成事件。
之后我使用信号量(semaphore)实现了如下模型:
def producer(RESP_q, URL_q, SEM): with SEM: while True: url = URL_q.get() if url == "END": break RESP = produce_txns(url) RESP_q.put(RESP) def consumer(RESP_q, SEM, NP): while SEM.get_value() < NP or not RESP_q.empty(): resp = RESP_q.get() for txn in resp: _txn = E_Transaction(txn) print(_txn) RESP_q.task_done() class Manager: def __init__(self): self.URL_q = Queue() self.RESP_q = JoinableQueue() self.max_processes = cpu_count() self.SEM = Semaphore(self.max_processes // 2) def start(self): self.worker = [] for i in range(0, self.max_processes, 2): self.worker.append(Process(target=producer, args=(self.RESP_q, self.URL_q, self.SEM))) self.worker.append(Process(target=consumer, args=(self.RESP_q, self.SEM, self.max_processes // 2))) url_server(self.URL_q, self.max_processes // 2) # Consider URL_q holds -> [*data, *["END"]*(self.max_processes // 2)] for worker in self.worker: worker.start() self.stop() def stop(self): for worker in self.worker: worker.join() self.RESP_q.join() self.RESP_q.close() self.URL_q.close() Manager().start()
存在的缺陷
该实现存在死锁风险:当消费者侧RESP_q为空、且SEM值接近设置的最大进程数时,如果解释器刚判断完while循环条件成立,SEM的值就变为等于最大进程数,此时所有生产者都已退出,程序就会阻塞在get方法处。
修复方案
编辑1:
@Louis Lac 的实现同样是正确的,我通过添加try-except块修改代码,解决了死锁问题:
def consumer(RESP_q, SEM, NP): while SEM.get_value() < NP or not RESP_q.empty(): try: resp = RESP_q.get(timeout=0.5) except Exception: continue
内容的提问来源于stack exchange,提问作者Rinkeby
相关产品推荐
相关产品推荐

