直接用 collections.deque 不行,因为它不是线程安全的,多线程调用 append() 或 popleft() 可能导致数据错乱或 indexerror,即使加锁也易在 len(q) 等边界判断时引发竞态。

为什么直接用 collections.deque 不行?
因为 deque 本身不是线程安全的——多个线程同时调用 append() 或 popleft() 可能导致数据错乱或 IndexError,即使只用一个锁包装它,也容易在边界条件(如满/空)下出现竞态:比如两个线程同时判断 len(q) 为真,然后都尝试 <code>append(),结果就超容。
用 threading.Semaphore 控制入队和出队的并发数
核心思路是把“容量控制”和“元素操作”拆开:用两个信号量分别表示“可用空位数”和“可用元素数”,避免检查-修改的竞态。实际代码里不需要手动维护长度或判断满/空,信号量会自然阻塞。
-
self._empty初始化为容量值,每次入队前acquire(),出队后release() -
self._full初始化为 0,每次出队前acquire(),入队后release() - 底层仍用
deque存储,但所有读写都在锁self._lock内进行,确保单次append()/popleft()原子性
from collections import deque
import threading
class ThreadSafeCircularQueue:
def __init__(self, maxsize):
self._queue = deque()
self._maxsize = maxsize
self._empty = threading.Semaphore(maxsize)
self._full = threading.Semaphore(0)
self._lock = threading.Lock()
def put(self, item):
self._empty.acquire()
with self._lock:
self._queue.append(item)
def get(self):
self._full.acquire()
with self._lock:
return self._queue.popleft()
注意 put() 和 get() 的阻塞行为与超时处理
默认情况下,put() 在队列满时会一直阻塞,get() 在空时也一样。这在多数服务场景中合理,但若需响应式退出(比如 worker 被中断),必须加 timeout 参数并捕获 threading.Semaphore 的超时异常——它抛的是 ValueError,不是 TimeoutError。
- 调用
self._empty.acquire(timeout=1)返回False表示超时,此时要主动release()已获取的资源(如果有) - 不要在
with self._lock:外部做耗时操作,否则锁持有时间不可控 - 如果业务需要非阻塞(如
put_nowait),应显式检查self._empty.acquire(blocking=False)返回值
为什么不用 queue.Queue?
queue.Queue 是线程安全的,但它内部用的是普通队列,不是循环结构——它的 maxsize 控制的是总长度,不会自动丢弃老元素;而“循环队列”隐含语义通常是固定容量、满时覆盖或拒绝。如果你真需要“满则覆盖”语义(比如日志缓冲),就得自己实现:在 put() 中检测长度,满时先 popleft() 再 append(),这时必须把整个操作包进同一把锁,且信号量逻辑也要相应调整(_full 不再简单 +1,得看是否真新增了元素)。
这个细节常被忽略:所谓“循环”,不单是数据结构形态,更是语义——是丢弃、阻塞,还是报错,得由你明确定义,然后选对应同步原语。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











