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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:33:56