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

Python多进程应用移除队列元素仍触发Full异常问题

问题分析

你的代码中出现Full异常的核心原因是多进程队列操作的竞态条件:当produce进程捕获到队列满的异常后,执行get_nowait()和后续put_nowait()的过程并非原子操作,期间listen进程的并发get_nowait()可能改变队列状态;同时listen进程的空循环(队列为空时直接重试)会占用大量CPU资源,加剧进程调度的不确定性,进一步增加竞态条件发生的概率。

解决方案

1. 用进程锁保证操作原子性

通过multiprocessing.Lock将"取旧元素+放新元素"的操作包裹为原子操作,避免并发操作导致的状态不一致:

import multiprocessing as mp
from queue import Full, Empty

def listen(queue: mp.Queue, lock: mp.Lock) -> None:
    while True:
        with lock:
            try:
                queue.get_nowait()
            except Empty:
                pass
            # 增加短暂休眠,减少空循环的CPU占用
            mp.sleep(0.001)

def produce(queue: mp.Queue, lock: mp.Lock) -> None:
    item = "qwerty"
    while True:
        try:
            queue.put_nowait(item)
        except Full:
            with lock:
                try:
                    # 原子操作:先取走一个旧元素,再放入新元素
                    queue.get_nowait()
                    queue.put_nowait(item)
                except Empty:
                    # 队列满时理论不会触发空异常,兜底处理
                    queue.put_nowait(item)
                except Full:
                    print("Full")

def run():
    queue = mp.Queue(maxsize=10)
    lock = mp.Lock()
    p1 = mp.Process(target=produce, args=(queue, lock))
    p2 = mp.Process(target=listen, args=(queue, lock))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

if __name__ == "__main__":
    run()

2. 优化listen进程的CPU占用

在listen进程的空循环中添加短暂休眠(mp.sleep(0.001)),减少无意义的CPU消耗,让produce进程的操作更稳定。

3. 自定义有界淘汰队列(可选)

如果需要更优雅的实现,可以封装一个自定义队列类,内置"满时淘汰旧元素"的逻辑,避免在业务代码中处理异常:

import multiprocessing as mp
from queue import Full, Empty

class BoundedReplaceQueue(mp.Queue):
    def __init__(self, maxsize: int):
        super().__init__(maxsize)
    
    def put_replace(self, item):
        try:
            self.put_nowait(item)
        except Full:
            # 原子操作:取旧放新
            with self._rlock:
                try:
                    self.get_nowait()
                except Empty:
                    pass
                self.put_nowait(item)

def listen(queue: BoundedReplaceQueue) -> None:
    while True:
        try:
            queue.get_nowait()
        except Empty:
            mp.sleep(0.001)

def produce(queue: BoundedReplaceQueue) -> None:
    item = "qwerty"
    while True:
        queue.put_replace(item)

def run():
    queue = BoundedReplaceQueue(maxsize=10)
    p1 = mp.Process(target=produce, args=(queue,))
    p2 = mp.Process(target=listen, args=(queue,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

if __name__ == "__main__":
    run()

为什么JoinableQueue没用?

JoinableQueue的核心作用是跟踪任务完成状态(通过task_done()和join()),用于等待所有任务被消费完成,无法解决队列满时的元素淘汰问题,因此不适用于你的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 12:50:57