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

Python多线程单生产者多消费者模式程序阻塞问题求助

问题修复方案

核心问题

  • 线程启动逻辑错误:循环内start()后直接调用join()会导致线程串行执行,先等生产者完全退出才启动第一个消费者,第一个消费者完全退出才启动第二个消费者,和多线程并发的设计目标完全不符。
  • 共享数据单元逻辑不支持多消费者:现有逻辑只要单个消费者读取完成就允许生产者写入新数据,没有等待所有消费者读取完毕,会导致后续消费者拿不到目标数据直接挂起。

修改后完整代码

import time, random
from threading import Thread, currentThread, Condition

class SharedCell(object):
    
    def __init__(self, consumer_num):
        self.data = -1
        self.writeable = True
        self.condition = Condition()
        self.consumer_num = consumer_num # 总消费者数量
        self.read_count = 0 # 已读取当前数据的消费者数量
        
    def setData(self, data):
        self.condition.acquire()
        while not self.writeable:
            self.condition.wait()
        
        print("%s setting data to %d" % \
              (currentThread().getName(), data))
        self.data = data
        self.writeable = False
        self.read_count = 0 # 新数据写入后重置读取计数
        self.condition.notifyAll()
        self.condition.release()
    
    def getData(self):
        self.condition.acquire()
        while self.writeable:
            self.condition.wait()
        print(f'accessing data {currentThread().getName()} {self.data}')
        self.read_count += 1
        # 所有消费者都读取完毕后才允许生产者写入
        if self.read_count == self.consumer_num:
            self.writeable = True
        self.condition.notifyAll()
        self.condition.release()
        return self.data
    
class Producer(Thread):
    
    def __init__(self, cell, accessCount, sleepMax):
        Thread.__init__(self, name = "Producer")
        self.accessCount = accessCount
        self.cell = cell
        self.sleepMax = sleepMax
    
    def run(self):
        print("%s starting up" % self.getName())
        for count in range(self.accessCount):
            time.sleep(random.randint(1, self.sleepMax))
            self.cell.setData(count + 1)
        print("%s is done producing\n" % self.getName())
        
class Consumer(Thread):
    
    def __init__(self, name, cell, accessCount, sleepMax):
        Thread.__init__(self, name=name)
        self.accessCount = accessCount
        self.cell = cell
        self.sleepMax = sleepMax
        
    def run(self):
        print("%s starting up" % self.getName())
        for count in range(self.accessCount):
            time.sleep(random.randint(1, self.sleepMax))
            value = self.cell.getData()
        print("%s is done consuming\n" % self.getName())
        
def main():
    accessCount = int(input("Enter the number of accesses: "))
    sleepMax = 4
    consumer_num = 2
    cell = SharedCell(consumer_num)
    
    producer = Producer(cell, accessCount, sleepMax)
    consumer1 = Consumer("Consumer1", cell, accessCount, sleepMax)
    consumer2 = Consumer("Consumer2", cell, accessCount, sleepMax)
    
    threads = [producer, consumer1, consumer2]
    
    print("Starting the threads")      
    # 先启动所有线程
    for thread in threads:
        thread.start()
    # 再统一等待所有线程执行完成
    for thread in threads:
        thread.join()
    print("All threads exit")

if __name__ == "__main__":
    main()

修改点说明

  1. 线程启动部分:拆分了启动和等待的逻辑,先把所有线程全部启动后再统一调用join(),实现真正的并发运行。
  2. SharedCell新增consumer_num参数记录消费者总数量,新增read_count记录当前数据已被读取的消费者次数。
  3. 消费者每次读取数据后计数+1,只有当计数等于总消费者数量时,才将writeable设为True,允许生产者写入新数据,符合「所有消费者均访问完已生产的数据后,生产者才可生产新数据」的规则。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 19:06:04