本文目录导读:

在Python中,队列安全(多线程/多进程安全)主要通过标准库提供的线程安全队列和进程安全队列来实现,以下是保障队列安全的核心案例和原理:
线程安全队列 - queue.Queue
基础案例
from queue import Queue
import threading
import time
# 创建线程安全队列
safe_queue = Queue(maxsize=10)
# 生产者
def producer():
for i in range(5):
item = f"数据-{i}"
safe_queue.put(item) # 线程安全操作
print(f"生产: {item}")
time.sleep(0.5)
# 消费者
def consumer():
while True:
item = safe_queue.get() # 线程安全操作
print(f"消费: {item}")
safe_queue.task_done()
time.sleep(1)
# 启动线程
t1 = threading.Thread(target=producer)
t2 = threading.Thread(target=consumer)
t1.start()
t2.start()
t1.join()
进程安全队列 - multiprocessing.Queue
跨进程案例
from multiprocessing import Process, Queue
def worker(q):
while True:
task = q.get()
if task is None: # 终止信号
break
print(f"进程{os.getpid()} 处理 {task}")
if __name__ == "__main__":
q = Queue()
# 启动多个工作进程
workers = []
for _ in range(3):
p = Process(target=worker, args=(q,))
p.start()
workers.append(p)
# 添加任务
for i in range(10):
q.put(f"任务-{i}")
# 发送终止信号(每个进程一个)
for _ in workers:
q.put(None)
for p in workers:
p.join()
保障队列安全的核心机制
① 内置锁机制
# Queue 内部实现(简化示意)
class SafeQueue:
def __init__(self):
self._mutex = threading.Lock() # 互斥锁
self._not_empty = threading.Condition(self._mutex)
self._not_full = threading.Condition(self._mutex)
def put(self, item, block=True, timeout=None):
with self._not_full: # 自动加锁
while self.full():
self._not_full.wait() # 队列满时等待
self._put(item)
self._not_empty.notify() # 通知消费者
def get(self, block=True, timeout=None):
with self._not_empty: # 自动加锁
while self.empty():
self._not_empty.wait() # 队列空时等待
item = self._get()
self._not_full.notify() # 通知生产者
return item
② 防止死锁设计
# 使用超时机制避免死锁
try:
item = safe_queue.get(timeout=5) # 5秒超时
except queue.Empty:
print("队列为空,超时退出")
# 非阻塞操作
if not safe_queue.empty():
item = safe_queue.get_nowait() # 不等待
高级安全队列模式
批量处理+异常处理
from queue import Queue, Empty
import threading
class SafeTaskQueue:
def __init__(self, maxsize=100):
self.queue = Queue(maxsize=maxsize)
self.running = threading.Event()
self.running.set()
def add_task(self, task, priority=0):
"""安全添加任务"""
self.queue.put((priority, task), timeout=5)
def get_batch(self, batch_size=10):
"""批量获取任务"""
tasks = []
for _ in range(batch_size):
try:
_, task = self.queue.get_nowait()
tasks.append(task)
except Empty:
break
return tasks
def stop(self):
"""安全停止"""
self.running.clear()
# 清空队列防止阻塞
while not self.queue.empty():
try:
self.queue.get_nowait()
except Empty:
break
# 使用示例
task_queue = SafeTaskQueue()
def worker():
while task_queue.running.is_set():
batch = task_queue.get_batch()
if batch:
for task in batch:
# 处理任务
pass
else:
threading.Event().wait(0.1)
# 启动工作线程
threads = [threading.Thread(target=worker) for _ in range(4)]
for t in threads:
t.daemon = True # 守护线程
t.start()
性能与安全的平衡
使用队列最佳实践
import queue
import time
# 1. 合理设置队列大小
q = queue.Queue(maxsize=1000) # 防止内存溢出
# 2. 使用 join() 确保所有任务完成
def worker():
while True:
item = q.get()
if item is None: # 哨兵值
break
# 处理任务
q.task_done()
# 等待所有任务完成
q.join()
# 3. 监控队列状态
q_size = q.qsize() # 当前队列大小
is_full = q.full()
is_empty = q.empty()
# 4. 优雅关闭
def shutdown_workers(workers, q):
"""安全关闭所有工作线程"""
for _ in workers:
q.put(None) # 发送终止信号
for w in workers:
w.join(timeout=5) # 超时等待
第三方增强队列
Redis队列(适合分布式场景)
import redis from rq import Queue as RQueue from rq.job import Job # Redis 支持的队列 r = redis.Redis() q = RQueue(connection=r) # 安全特性 job = q.enqueue(func, args=(1, 2)) result = job.result # 获取结果 status = job.get_status() # 检查状态
- 内置安全:
queue.Queue和multiprocessing.Queue已经内置了线程/进程安全机制 - 避免手动加锁:使用标准库队列通常比自己实现锁更安全
- 设置超时:防止永久阻塞
- 使用哨兵值:优雅地通知消费者停止
- 监控队列状态:避免队列溢出或空转
- 选择正确类别:单线程用
queue.Queue,多进程用multiprocessing.Queue,分布式用 Redis/RabbitMQ
这些案例展示了Python如何通过标准库和设计模式保障队列操作的安全性,避免了常见的竞态条件、死锁和数据不一致问题。