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

基于Python多进程与管道的生产者消费者问题实现求助(需满足三种运行场景)

基于Python多进程与管道的生产者消费者问题实现求助(需满足三种运行场景)

嘿,我仔细看了你的代码,发现几个核心问题导致程序没有输出甚至卡住,咱们一步步梳理和修复:

问题分析

  1. 多进程内存隔离导致共享变量失效:你直接定义的buffer、pro_counter、con_counter是每个进程各自的副本,生产者进程里添加的内容,消费者根本看不到,计数器也完全不同步,这是最致命的问题。
  2. Condition变量使用错误:你注释掉了condition.notify_all(),导致等待的进程永远无法被唤醒;而且Condition的wait逻辑和锁的配合也有问题,没有正确触发唤醒机制。
  3. 管道阻塞逻辑不合理:消费者一开始就调用input_pipe.recv(),如果生产者还没发送消息,会直接卡住;同时消费者里的while con_counter > 0循环会把计数器直接清零,完全不符合消费一个减一个的逻辑。
  4. 计数器维护逻辑混乱:生产者和消费者各自维护计数器,没有同步,导致缓冲区状态判断完全错误。

修复后的代码

下面是调整后的代码,我用multiprocessing.Manager创建了跨进程共享的缓冲区和计数器,修正了Condition的唤醒逻辑,同时保留了你要求的管道通信机制:

from multiprocessing import Process, Pipe, Lock, Condition, Manager
import time
import random

def create_items(output_pipe, notify_pipe, lock, condition, buffer, pro_counter, buffer_size):
    notify_conn = notify_pipe[0]
    for _ in range(100):
        item = random.randint(0, 1000)

        # 生产者等待缓冲区有空位
        with condition:
            while pro_counter.value >= buffer_size:
                print("Producer waiting: Buffer full.")
                condition.wait()

            # 添加物品到共享缓冲区
            with lock:
                buffer.append(item)
                pro_counter.value += 1
                print(f"Producer: Producing {item} | Buffer size: {len(buffer)}")
                output_pipe.send("produced")  # 通知消费者有新物品

            condition.notify_all()  # 唤醒等待的消费者

        # 处理消费者的消费确认
        while notify_conn.poll():
            ack = notify_conn.recv()
            if ack == "consumed":
                with lock:
                    pro_counter.value -= 1
                condition.notify_all()  # 唤醒等待的生产者

        time.sleep(0.01)  # 这里可以根据场景调整睡眠时间

    # 通知消费者生产结束
    output_pipe.send(None)
    output_pipe.close()
    notify_conn.close()

def consume_items(input_pipe, notify_pipe, lock, condition, buffer, con_counter, buffer_size):
    notify_conn = notify_pipe[1]
    while True:
        # 接收生产者的消息
        ack = input_pipe.recv()
        if ack is None:
            # 处理剩余物品
            while buffer:
                with lock:
                    item = buffer.pop(0)
                    con_counter.value -= 1
                    print(f"Consumer: Consumed {item} | Buffer size: {len(buffer)}")
                    notify_conn.send("consumed")
                time.sleep(1)  # 模拟消费延迟
            notify_conn.close()
            input_pipe.close()
            return
        elif ack == "produced":
            with lock:
                con_counter.value += 1

        # 消费者等待缓冲区有物品
        with condition:
            while con_counter.value <= 0:
                print("Consumer waiting: Buffer empty.")
                condition.wait()

            # 消费一个物品
            with lock:
                if buffer:
                    item = buffer.pop(0)
                    con_counter.value -= 1
                    print(f"Consumer: Consumed {item} | Buffer size: {len(buffer)}")
                    notify_conn.send("consumed")  # 通知生产者已消费

            condition.notify_all()  # 唤醒等待的生产者

        time.sleep(1)  # 这里可以根据场景调整睡眠时间

if __name__ == "__main__":
    buffer_size = 30
    # 使用Manager创建跨进程共享的变量
    manager = Manager()
    buffer = manager.list()
    pro_counter = manager.Value('i', 0)
    con_counter = manager.Value('i', 0)

    # 创建管道、锁和条件变量
    pipe_1 = Pipe(True)
    notify_pipe = Pipe(True)
    lock = Lock()
    condition = Condition(lock)

    # 创建进程(这里可以根据场景修改sleep时间)
    producer = Process(target=create_items, args=(pipe_1[1], notify_pipe, lock, condition, buffer, pro_counter, buffer_size))
    consumer = Process(target=consume_items, args=(pipe_1[0], notify_pipe, lock, condition, buffer, con_counter, buffer_size))

    consumer.start()
    producer.start()

    # 关闭主进程中不需要的管道端
    pipe_1[0].close()
    pipe_1[1].close()
    notify_pipe[0].close()
    notify_pipe[1].close()

    producer.join()
    consumer.join()

    print("Producer-Consumer processing complete.")

三种场景的设置方法

你只需要修改生产者和消费者函数里的time.sleep()参数即可:

  • 场景1:生产者和消费者同时工作:把生产者的sleep(0.01)改成sleep(1),消费者保持sleep(1),这样两者处理速度相近,会交替生产和消费。
  • 场景2:缓冲区满,生产者等待:生产者sleep(0.01),消费者sleep(1),生产者生产速度远快于消费,很快会填满缓冲区,触发生产者等待。
  • 场景3:缓冲区空,消费者等待:生产者sleep(1),消费者sleep(0.01),消费者消费速度远快于生产,很快会把缓冲区清空,触发消费者等待。

关键修复点说明

  1. 用Manager创建共享的缓冲区和计数器,确保跨进程能访问同一个状态。
  2. 修复了Condition的notify_all()调用,确保等待的进程能被正确唤醒。
  3. 调整了管道的接收逻辑,避免不必要的阻塞,同时处理了生产结束后的剩余物品。
  4. 修正了计数器的维护逻辑,确保生产者和消费者的计数器与实际缓冲区状态同步。

备注:内容来源于stack exchange,提问作者Phantom

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 12:37:59