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

Python多消费者线程需读取队列全部元素:Condition与ThreadPoolExecutor使用求助

多消费者线程共享读取队列元素的同步方案

问题背景

需要实现多个消费者线程读取队列中的每一个元素(即每个元素必须被所有消费者读取后才从队列移除),使用ThreadPoolExecutor提交任务。由于标准queue.get()会弹出元素,改用直接访问queue.queue[0]获取队首,但无法正确用threading.Condition在新元素可用时通知等待线程,需解决同步问题。

用户提供的示例代码:

def read():
    if (self.q.empty()) or (self.q.queue[0]==self.sentinel):
        raise Exception("Empty Queue")

    if (self not in self.read_ports):
        # Acquire lock
        self.condition.acquire()
        data = self.q.queue[0]
        self.read_ports.append(self)
        # Release lock
        self.condition.release()
        #print("Read ports1", self.read_ports)
    if self.is_identical(self.input_ports, self.read_ports):
        #print("Read ports2", self.read_ports)

        # All ports have read the top element, pop it now and clear read ports
        #Acquire lock
        self.q.get()
        self.read_ports.clear()
        # Release lock notify all waiting threads

解决方案

核心是用threading.Condition实现线程间的同步:当队列为空时让线程等待,当所有消费者读取完当前元素并弹出后,通知所有线程处理下一个元素;当新元素加入队列时,同样通知等待的线程。

完整实现代码

import threading
from concurrent.futures import ThreadPoolExecutor
import queue

class SharedQueue:
    def __init__(self, consumer_count, sentinel=None):
        self.q = queue.Queue()
        self.condition = threading.Condition()
        self.consumer_count = consumer_count
        self.read_ports = set()  # 用集合存线程标识,避免重复添加
        self.sentinel = sentinel

    def read(self):
        with self.condition:
            # 循环等待,直到队列非空且不是终止信号,或者遇到终止信号
            while self.q.empty() or (self.q.queue[0] == self.sentinel and not self.q.empty()):
                if not self.q.empty() and self.q.queue[0] == self.sentinel:
                    # 遇到终止信号,返回None让线程退出
                    return None
                self.condition.wait()

            # 获取当前队首元素
            current_data = self.q.queue[0]
            # 记录当前线程已读取该元素(用线程ID作为标识)
            thread_id = threading.get_ident()
            self.read_ports.add(thread_id)

            # 检查是否所有消费者都读取了当前元素
            if len(self.read_ports) == self.consumer_count:
                # 弹出元素,清空已读记录
                self.q.get()
                self.read_ports.clear()
                # 通知所有等待的线程,准备处理下一个元素
                self.condition.notify_all()

            return current_data

    def put(self, item):
        with self.condition:
            self.q.put(item)
            # 新元素加入后,通知所有等待的线程
            self.condition.notify_all()

def consumer(shared_queue, consumer_name):
    while True:
        data = shared_queue.read()
        if data is None:
            print(f"消费者 {consumer_name} 收到终止信号,退出")
            break
        print(f"消费者 {consumer_name} 读取到元素: {data}")

if __name__ == "__main__":
    # 3个消费者线程,3个队列元素
    CONSUMER_NUM = 3
    shared_queue = SharedQueue(CONSUMER_NUM, sentinel="STOP")

    # 添加测试元素
    for i in range(3):
        shared_queue.put(f"元素{i+1}")
    # 添加终止信号
    shared_queue.put("STOP")

    # 用ThreadPoolExecutor启动消费者
    with ThreadPoolExecutor(max_workers=CONSUMER_NUM) as executor:
        for i in range(CONSUMER_NUM):
            executor.submit(consumer, shared_queue, f"Thread-{i+1}")

关键说明

  1. Condition的正确使用:所有访问队列、read_ports的操作都在with self.condition上下文(自动加锁/解锁)中进行,避免竞态条件。wait()会释放锁并阻塞线程,直到被notify_all()唤醒。
  2. 线程标识:用threading.get_ident()获取线程ID作为已读记录的标识,避免用线程对象本身导致的判断问题。
  3. 终止信号处理:当队列遇到sentinel时,让所有线程退出循环。
  4. ThreadPoolExecutor配合:线程池中的线程本质是普通的threading.Thread实例,直接使用共享的Condition对象即可,无需特殊适配——只要确保所有线程访问的是同一个SharedQueue实例。

内容的提问来源于stack exchange,提问作者Lokesh jain

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 02:25:51