Python的queue.Queue阻塞时,get/put方法是否持有队列锁?
Python多线程生产者/消费者模式实现(基于queue.Queue)
用Python搞多线程生产者/消费者,queue.Queue真的是最省心的选择——它内置了线程安全的同步机制,不用咱们手动加锁处理竞态问题,直接拿来用就行。下面我把核心实现逻辑,尤其是消费者的两种消费方式给你讲清楚:
核心架构说明
生产者线程负责生成数据,把数据塞进queue.Queue队列;消费者线程从队列里取数据处理。队列作为中间缓冲区,自动平衡生产和消费的速度差。
消费者的两种消费方式
1. 阻塞式消费
这种方式是最常用的,当self.running为True时,消费者会一直阻塞等待队列中有数据,直到拿到数据再处理:
while self.running: # 阻塞等待队列出现数据 data = self.receiver_q.get(block=True) # 这里写你的数据处理逻辑 process_data(data) # 记得告诉队列任务完成(如果用了join()机制的话) self.receiver_q.task_done()
特点&注意事项:
- 优点:不用轮询,空闲时不占用CPU资源,适合数据产生频率不固定的场景。
- 优雅停止技巧:如果要让消费者退出,除了把
self.running设为False,最好往队列里塞一个“终止信号”(比如None),让阻塞的线程拿到信号后主动退出循环,避免线程一直卡着。
2. 非阻塞式消费
如果你的消费者需要在等待数据的间隙做其他事情(比如定期检查状态、处理其他小任务),可以用非阻塞模式:
import queue while self.running: try: # 非阻塞获取数据,队列空就抛出异常 data = self.receiver_q.get(block=False) # 处理数据 process_data(data) self.receiver_q.task_done() except queue.Empty: # 队列空的时候,可以在这里做其他事情 do_other_thing() # 也可以加个短暂休眠,减少CPU轮询消耗 time.sleep(0.1)
特点&注意事项:
- 优点:灵活性高,等待间隙能处理其他逻辑。
- 缺点:会频繁轮询队列,比阻塞模式消耗更多CPU,建议加个
time.sleep()来降低轮询频率。
完整示例参考
给你写个简单的可运行示例,包含生产者和两种消费者:
import threading import queue import time import random class Producer(threading.Thread): def __init__(self, q): super().__init__() self.q = q self.running = True def run(self): while self.running: # 生成随机数据 data = random.randint(1, 100) self.q.put(data) print(f"生产者生成数据: {data}") time.sleep(random.uniform(0.5, 1.5)) class BlockingConsumer(threading.Thread): def __init__(self, q): super().__init__() self.q = q self.running = True def run(self): while self.running: data = self.q.get(block=True) if data is None: # 终止信号 break print(f"阻塞消费者处理数据: {data}") self.q.task_done() time.sleep(0.3) class NonBlockingConsumer(threading.Thread): def __init__(self, q): super().__init__() self.q = q self.running = True def run(self): while self.running: try: data = self.q.get(block=False) if data is None: # 终止信号 break print(f"非阻塞消费者处理数据: {data}") self.q.task_done() time.sleep(0.3) except queue.Empty: print("非阻塞消费者等待中,做点别的...") time.sleep(0.5) if __name__ == "__main__": q = queue.Queue(maxsize=5) producer = Producer(q) consumer1 = BlockingConsumer(q) consumer2 = NonBlockingConsumer(q) producer.start() consumer1.start() consumer2.start() # 运行5秒后停止 time.sleep(5) producer.running = False # 发送终止信号给消费者 q.put(None) q.put(None) producer.join() consumer1.join() consumer2.join() print("所有线程结束")
内容的提问来源于stack exchange,提问作者wrkyle
相关产品推荐
相关产品推荐

