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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 04:55:07