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

Python线程:生产者-消费者模式中如何正确保障队列顺序?

生产者-多Worker-有序输出队列的实现优化问题

我有一个定期生成数据的生产者,数据由多个Worker通过网络分析,要求结果严格按照生产者输出的顺序进行后续处理,而非Worker完成分析的顺序。架构如下:
producer -> queue1 -> worker1, ..., workern -> queue2 -> ...

我想了解当前实现是否符合Python风格,是否存在更优方案。我的方案核心是:Worker处理完数据后,等待合适时机再将结果放入最终队列。以下是代码实现:

from threading import Thread, Lock
from queue import Queue
from time import sleep
from random import random

class Producer(Thread):

    def __init__(self, queue, id=1):
        super().__init__()
        self.queue = queue
        self.id = id

    def run(self, init=1, end=20):
        for n in range(init, end+1):
            item = f"item #{n}"
            self.queue.put(item)
            print(f"producer #{self.id}: {item} sent to queue")
            sleep(1)
        self.queue.put(None)

class Worker(Thread):
    lock1 = Lock()
    lock2 = Lock()

    def __init__(self, queue1, queue2, id=1):
        super().__init__()
        self.queue1 = queue1
        self.queue2 = queue2
        self.id = id

    def run(self):
        q1, q2 = self.queue1, self.queue2
        cons = f"consumer #{self.id}"

        while True:
            with self.lock1:
                item = q1.get()
                if item is None:
                    q1.put(None)
                    return
                q1.nread += 1
                n = q1.nread

            print(f"{cons}: {item} obtained [get #{n}]")
            sleep(6 * random())
            print(f"{cons}: {item} processed")

            while True:
                with self.lock2:
                    if n == q2.nwritten + 1:
                        q2.put(item)
                        q2.nwritten += 1
                        print(" " * 10 + f"-> {cons}: {item} sent to queue")
                        q1.task_done()
                        break
                sleep(.1)

q1 = Queue()
q1.nread = 0
q2 = Queue()
q2.nwritten = 0

workers = [Worker(q1, q2, id=i) for i in range(1, 5)]
producer = Producer(q1)

for worker in workers:
    worker.start()

producer.start()
producer.join()
for worker in workers:
    worker.join()

现有实现的优缺点分析

不符合Python风格的地方

  1. 直接修改标准库类实例属性:给Queue实例添加nread和nwritten属性,打破了封装原则。标准库Queue未公开这些属性,此类“猴子补丁”式修改容易引发未知问题,违背Python“显式优于隐式”的设计原则。
  2. 忙等浪费CPU资源:Worker处理完数据后,通过while True循环加sleep(0.1)等待写入时机,属于低效的忙等模式,会不必要消耗CPU。
  3. 锁的设计不够清晰:使用Worker类的类级锁lock1和lock2,虽然能保证计数线程安全,但锁的作用范围和职责未明确区分,可读性差。
  4. 终止信号处理不严谨:用None作为终止信号,多个Worker会重复将None放回队列,逻辑冗余。

可取之处

核心思路正确:通过追踪任务序号,确保只有当前应输出的任务结果被放入queue2,严格保证了输出顺序。


更优的Python风格实现方案

推荐采用带序号的任务+结果字典+有序输出线程的模式,或直接使用concurrent.futures.ThreadPoolExecutor简化线程管理,同时保证结果有序。以下是两种实现方式:

方案1:基于ThreadPoolExecutor的简化实现

利用ThreadPoolExecutor管理Worker线程,通过字典收集结果后按序号输出到queue2,代码简洁且符合标准库使用习惯:

from concurrent.futures import ThreadPoolExecutor
from queue import Queue
from time import sleep
from random import random
from threading import Thread

def producer(queue, init=1, end=20):
    for n in range(init, end+1):
        item = (n, f"item #{n}")
        queue.put(item)
        print(f"producer: {item[1]} sent to queue")
        sleep(1)
    # 发送结束信号
    queue.put((None, None))

def worker_task(item):
    idx, data = item
    worker_id = idx % 4 + 1
    print(f"worker {worker_id}: {data} obtained")
    sleep(6 * random())
    print(f"worker {worker_id}: {data} processed")
    return (idx, data)

def ordered_output(queue_in, queue_out):
    results = {}
    next_idx = 1
    while True:
        task = queue_in.get()
        if task[0] is None:
            break
        future = executor.submit(worker_task, task)
        results[task[0]] = future

        # 检查并输出连续可用的结果
        while next_idx in results and results[next_idx].done():
            idx, data = results[next_idx].result()
            queue_out.put(data)
            print(f"    -> {data} sent to queue2")
            del results[next_idx]
            next_idx += 1
        queue_in.task_done()

# 初始化队列
queue1 = Queue()
queue2 = Queue()

# 启动线程池和输出线程
executor = ThreadPoolExecutor(max_workers=4)
output_thread = Thread(target=ordered_output, args=(queue1, queue2))
output_thread.start()

# 启动生产者
producer(queue1)
queue1.join()
executor.shutdown()
output_thread.join()

方案2:自定义线程的优化版本

若坚持使用自定义Thread类,优化点包括:用线程安全的计数器替代直接修改Queue属性,用threading.Condition替代忙等,明确锁的职责:

from threading import Thread, Lock, Condition
from queue import Queue
from time import sleep
from random import random

class Producer(Thread):
    def __init__(self, queue, id=1):
        super().__init__()
        self.queue = queue
        self.id = id

    def run(self, init=1, end=20):
        for n in range(init, end+1):
            item = (n, f"item #{n}")
            self.queue.put(item)
            print(f"producer #{self.id}: {item[1]} sent to queue")
            sleep(1)
        # 发送结束信号
        self.queue.put((None, None))

class Worker(Thread):
    def __init__(self, queue_in, result_dict, output_condition, id=1):
        super().__init__()
        self.queue_in = queue_in
        self.result_dict = result_dict
        self.output_condition = output_condition
        self.id = id

    def run(self):
        worker_name = f"consumer #{self.id}"
        while True:
            item = self.queue_in.get()
            idx, data = item
            if idx is None:
                self.queue_in.put(item)
                break

            print(f"{worker_name}: {data} obtained [task #{idx}]")
            sleep(6 * random())
            print(f"{worker_name}: {data} processed")

            # 存储结果并通知输出线程
            with self.output_condition:
                self.result_dict[idx] = data
                self.output_condition.notify()
            self.queue_in.task_done()

class OrderedOutputThread(Thread):
    def __init__(self, queue_out, result_dict, output_condition, total_tasks=20):
        super().__init__()
        self.queue_out = queue_out
        self.result_dict = result_dict
        self.output_condition = output_condition
        self.next_idx = 1
        self.total_tasks = total_tasks

    def run(self):
        while self.next_idx <= self.total_tasks:
            with self.output_condition:
                # 等待当前序号的结果生成
                while self.next_idx not in self.result_dict:
                    self.output_condition.wait()
                # 输出结果到queue2
                data = self.result_dict[self.next_idx]
                self.queue_out.put(data)
                print(f"    -> {data} sent to queue")
                del self.result_dict[self.next_idx]
                self.next_idx += 1

# 初始化共享资源
queue1 = Queue()
queue2 = Queue()
result_dict = {}
output_condition = Condition()

# 启动线程
workers = [Worker(queue1, result_dict, output_condition, id=i) for i in range(1,5)]
producer = Producer(queue1)
output_thread = OrderedOutputThread(queue2, result_dict, output_condition)

output_thread.start()
for worker in workers:
    worker.start()
producer.start()

producer.join()
queue1.join()
for worker in workers:
    worker.join()
output_thread.join()

方案优势对比

  1. 封装性更好:不再修改标准库类属性,用专用共享对象(如result_dict、Condition)管理状态,符合Python面向对象原则。
  2. 避免忙等:使用threading.Condition的wait()/notify()机制,仅当有结果可用时唤醒输出线程,节省CPU资源。
  3. 代码可读性更高:生产者、Worker、输出线程职责划分清晰,逻辑直观。
  4. 终止处理更严谨:明确的结束信号传递,避免重复放入终止标记。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 22:30:59