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风格的地方
- 直接修改标准库类实例属性:给
Queue实例添加nread和nwritten属性,打破了封装原则。标准库Queue未公开这些属性,此类“猴子补丁”式修改容易引发未知问题,违背Python“显式优于隐式”的设计原则。 - 忙等浪费CPU资源:Worker处理完数据后,通过
while True循环加sleep(0.1)等待写入时机,属于低效的忙等模式,会不必要消耗CPU。 - 锁的设计不够清晰:使用Worker类的类级锁
lock1和lock2,虽然能保证计数线程安全,但锁的作用范围和职责未明确区分,可读性差。 - 终止信号处理不严谨:用
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()
方案优势对比
- 封装性更好:不再修改标准库类属性,用专用共享对象(如
result_dict、Condition)管理状态,符合Python面向对象原则。 - 避免忙等:使用
threading.Condition的wait()/notify()机制,仅当有结果可用时唤醒输出线程,节省CPU资源。 - 代码可读性更高:生产者、Worker、输出线程职责划分清晰,逻辑直观。
- 终止处理更严谨:明确的结束信号传递,避免重复放入终止标记。
内容的提问来源于stack exchange,提问作者loqueelviento
相关产品推荐
相关产品推荐

