要求: - 任意数量的生产者/消费者并发调用都必须正确;close 可以与 push/pop 并发调用。 - 阻塞必须是真正的挂起等待,不允许 sleep/yield 轮询或自旋忙等。 - 只使用标准库,底层容器可以直接用 std::queue或用 C + pthread 实现,接口语义等价即可。 - 如果你认为接口签名有更合理的设计,可以修改,但请说明理由。
所有共享状态只由同一把互斥锁保护,并使用两个条件变量管理“队列非空”和“队列未满”。
# deque 是双端队列,append / popleft 都是高效操作。
from collections import deque
# Lock 提供互斥访问;Condition 用于“条件不满足时挂起等待”。
from threading import Condition, Lock
# 用于泛型类型标注。
from typing import Deque, Generic, Optional, TypeVar
# 用于实现可靠的超时等待。
import time
# T 表示队列中元素的任意类型。
T = TypeVar("T")
# Generic[T] 表示 BoundedBlockingQueue 可以保存任意类型的元素。
class BoundedBlockingQueue(Generic[T]):
# capacity 是队列能保存的最大元素数量。
def __init__(self, capacity: int):
# 容量为 0 时,任何 push 永远无法成功,通常应视为非法参数。
if capacity <= 0:
raise ValueError("capacity 必须大于 0")
# 保存容量上限;它创建后不再改变。
self._capacity = capacity
# 实际存放数据的容器。
self._queue: Deque[T] = deque()
# 表示队列是否已关闭;关闭后不再接受新的 push。
self._closed = False
# 这一把锁保护所有共享状态:_queue、_closed、_capacity。
# 不能让生产和消费分别使用不同的锁,否则会发生竞态条件。
self._mutex = Lock()
# 当队列为空时,消费者在这个条件变量上等待。
# 它与 _not_full 共用同一把 _mutex。
self._not_empty = Condition(self._mutex)
# 当队列已满时,生产者在这个条件变量上等待。
# 两个 Condition 必须共享同一把锁,才能安全检查和修改队列状态。
self._not_full = Condition(self._mutex)
# 放入一个元素。
# 成功返回 True;若队列已关闭,返回 False。
def push(self, item: T) -> bool:
# 进入临界区并自动获取 _mutex。
# 离开 with 代码块时,无论正常返回还是异常,锁都会自动释放。
with self._not_full:
# 队列满且尚未关闭时,生产者必须等待。
# 必须用 while,不能用 if:
# 1. wait 可能发生“虚假唤醒”;
# 2. 被唤醒后,其他生产者可能已先抢到空位;
# 3. close() 也可能唤醒当前线程。
while len(self._queue) >= self._capacity and not self._closed:
# wait() 会先释放 _mutex,再挂起当前线程。
# 被 notify 唤醒后,wait() 会重新获得 _mutex 才返回。
self._not_full.wait()
# 退出 while 后,可能是队列有空位,也可能是队列关闭了。
# 关闭优先:关闭后绝不能再放入新元素。
if self._closed:
return False
# 此时持有锁、队列未满,因此可以安全加入元素。
self._queue.append(item)
# 队列从“可能为空”变成“非空”,唤醒一个等待消费的线程。
# notify() 唤醒一个即可,因为一次 push 最多满足一个 pop。
self._not_empty.notify()
# 表示元素已成功入队。
return True
# 取出一个元素。
# 队列关闭且所有遗留元素都取完时,返回 None。
def pop(self) -> Optional[T]:
# 获取同一把 _mutex,保证与 push()/close() 的操作互斥。
with self._not_empty:
# 队列为空、且尚未关闭时,消费者需要等待生产者放入元素。
while not self._queue and not self._closed:
# 临时释放锁并进入真正的阻塞等待,不会忙等或占用 CPU。
self._not_empty.wait()
# 能离开循环有两种可能:
# 1. 队列中有元素;
# 2. 队列已经关闭。
#
# 若仍为空,则只能是“关闭且元素已全部消费”。
if not self._queue:
return None
# 此时队列必定非空,安全地取出最早放入的元素。
item = self._queue.popleft()
# 队列从“可能已满”变成“至少有一个空位”,
# 唤醒一个等待 push 的生产者。
self._not_full.notify()
# 返回取出的元素。
return item
# 最多等待 timeout 秒取出一个元素。
# 成功返回 (True, item);
# 超时或关闭且队列已空时,返回 (False, None)。
def pop_with_timeout(self, timeout: float) -> tuple[bool, Optional[T]]:
# monotonic() 是单调递增时钟,不受系统时间被修改影响。
# 所以比 time.time() 更适合计算超时。
deadline = time.monotonic() + timeout
# 获取保护队列状态的同一把锁。
with self._not_empty:
# 空且未关闭时,继续等待;循环仍是为了防虚假唤醒和竞争。
while not self._queue and not self._closed:
# 计算本次还可等待多少秒。
remaining = deadline - time.monotonic()
# 剩余时间已用尽,说明超时。
if remaining <= 0:
return False, None
# 最多等待 remaining 秒。
# 生产者 push、close 唤醒、或等待超时后都会返回并重新检查条件。
self._not_empty.wait(remaining)
# 被唤醒后仍没有元素,说明队列已关闭且已经取空。
if not self._queue:
return False, None
# 队列非空,取出一个元素。
item = self._queue.popleft()
# 释放了一个容量位置,通知一个可能在等待的生产者。
self._not_full.notify()
# 返回成功标志和取到的元素。
return True, item
# 关闭队列。
# 已进入队列的元素仍可被 pop() 取出。
def close(self) -> None:
# 与 push/pop 竞争同一把锁,保证关闭状态改变是原子的。
with self._mutex:
# close() 设计为幂等操作:多次调用结果相同。
if self._closed:
return
# 标记为关闭;之后所有 push 都会失败。
self._closed = True
# 唤醒所有等待 pop 的消费者。
# 它们醒来后:
# - 若有遗留元素,继续消费;
# - 若队列为空,返回 None / False。
self._not_empty.notify_all()
# 唤醒所有等待 push 的生产者。
# 它们醒来后发现 _closed 为 True,返回 False。
self._not_full.notify_all()
你在面试中应该能把它的逻辑概括为:
生产者:
队列满 -> 等待 not_full
队列关闭 -> 失败
放入元素 -> 通知 not_empty
消费者:
队列空且未关闭 -> 等待 not_empty
队列关闭且为空 -> 结束
取出元素 -> 通知 not_full
关闭:
标记 closed
唤醒所有等待中的生产者和消费者
尤其要避免原题代码中的几个典型问题:
lock_guard不能用于condition_variable.wait();等待时必须使用可解锁、可重新加锁的锁对象(C++ 是std::unique_lock)。- 生产和消费不能各用一把锁;它们操作的是同一个队列,必须共享同一把锁。
capacity_应表示固定上限,不能随push/pop增减;当前元素数量由queue_.size()表示。close()必须唤醒所有等待线程,否则满队列上的生产者或空队列上的消费者可能永远无法返回。pop()应从队首取元素;Python 是popleft(),C++ 是front()后pop()。- 超时不是“先等普通条件变量、再
try_lock_for”;而是对条件变量本身进行限时等待。