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()
修改点说明
- 线程启动部分:拆分了启动和等待的逻辑,先把所有线程全部启动后再统一调用
join(),实现真正的并发运行。 - SharedCell新增
consumer_num参数记录消费者总数量,新增read_count记录当前数据已被读取的消费者次数。 - 消费者每次读取数据后计数+1,只有当计数等于总消费者数量时,才将
writeable设为True,允许生产者写入新数据,符合「所有消费者均访问完已生产的数据后,生产者才可生产新数据」的规则。
内容的提问来源于stack exchange,提问作者Anther Eros
相关产品推荐
相关产品推荐

