Python中实现Consumer完成任务后通知Producer的最优方案
需求与问题
我要并行处理多个文本文件的读取和分析流程,设置了一个Producer(生产者)和两个Consumer(消费者)——TextAnalyzer A和TextAnalyzer B。Producer维护一个RequestQueue,里面存放分配给两个消费者的任务。当TextAnalyzer A或B完成任务后,需要通知Producer;等两个消费者都完成各自任务后,Producer要启动ResultSender线程执行后续操作。
问题:在Python中,让Producer感知Consumer完成工作的最优方式有哪些?
现有实现代码
from time import sleep from random import random from threading import Thread from queue import Queue # 全局队列供生产者消费者通信 RequestQueue = Queue() FinishedQueue = Queue() class Task: def __init__(self, name, index, value, done): self.name = name self.index = index self.value = value self.done = done def getName(self): return self.name def getIndex(self): return self.index def getValue(self): return self.value def getDone(self): return self.done def setDone(self, done): self.done = done def producer(): print('Producer: Running') # 生成任务 task = Task("(csv)", 1, "test", False) RequestQueue.put(task) print(">>producer added {} {} {} {}".format( task.getName(), task.getIndex(), task.getValue(), task.getDone())) value = random() sleep(value) while True: if not FinishedQueue.empty(): fq = FinishedQueue.get() print(">>producer found {} {} {} {}".format( fq.getName(), fq.getIndex(), fq.getValue(), fq.getDone())) break sleep(1) RequestQueue.put(None) print('Producer: Done') # 消费者任务 def consumer(): print('Consumer: Running') while True: item = RequestQueue.get() # 检查停止信号 if item is None: break print(">>consumer is working on {} {} {} {}".format( item.getName(), item.getIndex(), item.getValue(), item.getDone())) # 模拟工作 sleep(1) # 标记任务完成 item.setDone(True) print(">>consumer finished {} {} {} {}".format( item.getName(), item.getIndex(), item.getValue(), item.getDone())) FinishedQueue.put(item) # 任务完成 print('Consumer: Done') def test(): tc = Thread(target=consumer, args=()) tc.start() tp = Thread(target=producer, args=()) tp.start() tp.join() tc.join() test()
几种最优实现方式
1. 队列+任务计数(现有思路优化)
基于你现有的队列通信思路,优化后可以更优雅地跟踪任务完成状态:
- 避免全局队列,将队列作为参数传递给线程,降低耦合
- 用队列自带的
task_done()和join()方法替代轮询 - 给任务标记归属的消费者,统计每个消费者的完成任务数,判断是否全部完成
示例优化代码:
from time import sleep from random import random from threading import Thread from queue import Queue class Task: def __init__(self, name, index, value, consumer_id): self.name = name self.index = index self.value = value self.consumer_id = consumer_id # 标记任务归属的消费者 def producer(): print('Producer: Running') request_queue = Queue() finished_queue = Queue() total_tasks_per_consumer = 2 # 每个消费者分配2个任务 # 生成并分配任务 for i in range(total_tasks_per_consumer): request_queue.put(Task("(csv)", i+1, f"task_{i+1}", "A")) request_queue.put(Task("(csv)", i+3, f"task_{i+3}", "B")) # 启动两个消费者 consumer_a = Thread(target=consumer, args=(request_queue, finished_queue, "A")) consumer_b = Thread(target=consumer, args=(request_queue, finished_queue, "B")) consumer_a.start() consumer_b.start() # 跟踪消费者完成状态 completed_count = {"A": 0, "B": 0} while True: if completed_count["A"] == total_tasks_per_consumer and completed_count["B"] == total_tasks_per_consumer: break finished_task = finished_queue.get() completed_count[finished_task.consumer_id] += 1 print(f">>Producer 收到 {finished_task.consumer_id} 完成的任务: {finished_task.name} {finished_task.index}") # 发送停止信号 request_queue.put(None) request_queue.put(None) consumer_a.join() consumer_b.join() # 启动ResultSender result_sender = Thread(target=result_sender) result_sender.start() result_sender.join() print('Producer: Done') def consumer(request_queue, finished_queue, consumer_id): print(f'Consumer {consumer_id}: Running') while True: item = request_queue.get() if item is None: break print(f">>Consumer {consumer_id} 处理任务: {item.name} {item.index} {item.value}") # 模拟分析工作 sleep(random()) finished_queue.put(item) request_queue.task_done() print(f'Consumer {consumer_id}: Done') def result_sender(): print('ResultSender: 开始执行后续操作') sleep(1) print('ResultSender: 完成') if __name__ == "__main__": producer()
2. 用threading.Event标记消费者完成
如果不需要跟踪单个任务,只需要知道消费者是否全部完成,用Event对象更简洁:给每个消费者分配一个Event,消费者完成所有任务后触发该Event,Producer等待两个Event都触发后启动ResultSender。
示例代码:
from time import sleep from random import random from threading import Thread, Event from queue import Queue class Task: def __init__(self, name, index, value): self.name = name self.index = index self.value = value def producer(): print('Producer: Running') request_queue = Queue() # 创建两个消费者的完成事件 a_done = Event() b_done = Event() # 生成任务 for i in range(4): request_queue.put(Task("(csv)", i+1, f"task_{i+1}")) # 启动消费者 consumer_a = Thread(target=consumer, args=(request_queue, a_done, "A")) consumer_b = Thread(target=consumer, args=(request_queue, b_done, "B")) consumer_a.start() consumer_b.start() # 等待两个消费者完成 a_done.wait() b_done.wait() print('Producer: 两个消费者都已完成') # 启动ResultSender result_sender = Thread(target=result_sender) result_sender.start() result_sender.join() print('Producer: Done') def consumer(request_queue, done_event, consumer_id): print(f'Consumer {consumer_id}: Running') while True: try: # 非阻塞获取任务,队列空则退出 item = request_queue.get(block=False) except: break print(f">>Consumer {consumer_id} 处理任务: {item.name} {item.index} {item.value}") sleep(random()) request_queue.task_done() done_event.set() # 触发完成事件 print(f'Consumer {consumer_id}: Done') def result_sender(): print('ResultSender: 开始执行后续操作') sleep(1) print('ResultSender: 完成') if __name__ == "__main__": producer()
3. 用concurrent.futures.ThreadPoolExecutor简化流程
如果不想手动管理线程,ThreadPoolExecutor可以大幅简化代码,它自带任务跟踪和结果回收功能,适合快速实现并行任务:
示例代码:
from time import sleep from random import random from concurrent.futures import ThreadPoolExecutor class Task: def __init__(self, name, index, value): self.name = name self.index = index self.value = value def process_task(task, consumer_id): print(f">>Consumer {consumer_id} 处理任务: {task.name} {task.index} {task.value}") sleep(random()) print(f">>Consumer {consumer_id} 完成任务: {task.name} {task.index}") return task def result_sender(): print('ResultSender: 开始执行后续操作') sleep(1) print('ResultSender: 完成') def producer(): print('Producer: Running') # 生成任务列表 tasks = [Task("(csv)", i+1, f"task_{i+1}") for i in range(4)] # 创建线程池,指定2个线程对应两个消费者 with ThreadPoolExecutor(max_workers=2) as executor: futures = [] # 分配任务给两个消费者 for idx, task in enumerate(tasks): consumer_id = "A" if idx % 2 == 0 else "B" futures.append(executor.submit(process_task, task, consumer_id)) # 等待所有任务完成 for future in futures: future.result() print('Producer: 所有消费者任务完成') # 启动ResultSender result_sender() print('Producer: Done') if __name__ == "__main__": producer()
总结
- 需跟踪单个任务完成情况:队列+任务计数的方式最灵活,能精准掌握每个任务的状态
- 仅需判断消费者是否全部完成:Event对象实现简单,代码耦合度低
- 不想手动管理线程:ThreadPoolExecutor简化流程,适合快速开发并行任务场景
内容的提问来源于stack exchange,提问作者takanoha
相关产品推荐
相关产品推荐

