Python多进程单写多读:并发读取与写优先实现问询
多智能体读写进程实现(Writer优先)
问题分析
你需要实现Writer优先的读写锁机制,核心需求:
- 空闲状态下,多个Reader可同时读取共享资源
- Writer需要写入时,拥有最高访问优先级,独占资源,且等待中的Writer会优先于后续Reader获取权限
原代码存在以下问题:
- 单一Condition变量无法实现多Reader并发,因为Condition的with块是互斥的,一次仅允许一个进程进入
- Reader中的
wait()无条件判断,会导致无意义阻塞 - 进程创建时target指向错误,未绑定类实例的方法
- 添加Writer进程时变量
i未定义
修正后的实现代码
import multiprocessing import time class RWLock: def __init__(self): self.condition = multiprocessing.Condition() self.reader_count = 0 # 当前活跃的Reader数量 self.writer_waiting = 0 # 等待中的Writer数量 self.is_writing = False # 是否有Writer正在写入 def acquire_read(self): with self.condition: # 有Writer在写或等待时,Reader阻塞 while self.is_writing or self.writer_waiting > 0: self.condition.wait() self.reader_count += 1 def release_read(self): with self.condition: self.reader_count -= 1 # 无活跃Reader时,通知等待的Writer if self.reader_count == 0: self.condition.notify_all() def acquire_write(self): with self.condition: self.writer_waiting += 1 # 等待所有Reader结束、无其他Writer写入 while self.is_writing or self.reader_count > 0: self.condition.wait() self.writer_waiting -= 1 self.is_writing = True def release_write(self): with self.condition: self.is_writing = False # 优先通知Writer保证优先级,再通知Reader self.condition.notify_all() class SharedData: def __init__(self): self.value = 0 self.rw_lock = RWLock() def modify_value(self, agent_id): self.rw_lock.acquire_write() try: print(f"Writer {agent_id} 开始写入,当前值: {self.value}") time.sleep(0.1) # 模拟写入耗时 self.value += 1 print(f"Writer {agent_id} 写入完成,新值: {self.value}") finally: self.rw_lock.release_write() def read_value(self, agent_id): self.rw_lock.acquire_read() try: print(f"Reader {agent_id} 读取值: {self.value}") time.sleep(0.05) # 模拟读取耗时 finally: self.rw_lock.release_read() if __name__ == "__main__": shared_data = SharedData() # 创建4个Reader进程和1个Writer进程 agents = [] # 添加Reader进程 for i in range(4): agents.append(multiprocessing.Process(target=shared_data.read_value, args=(i,))) # 添加Writer进程 agents.append(multiprocessing.Process(target=shared_data.modify_value, args=(0,))) # 启动所有进程 for agent in agents: agent.start() # 等待所有进程结束 for agent in agents: agent.join() print(f"最终值: {shared_data.value}")
关键实现说明
- RWLock类:自定义读写锁,通过三个状态变量控制访问逻辑:
reader_count:记录当前活跃Reader数量,实现多Reader并发读取writer_waiting:记录等待的Writer数量,保证Writer优先级(有Writer等待时,新Reader无法获取锁)is_writing:标记是否有Writer正在写入,确保写入操作独占资源
- Writer优先级保障:当有Writer等待时,后续调用
acquire_read()的Reader会被阻塞,直到所有Writer完成操作 - 异常安全:用
try-finally块确保锁的释放,避免进程异常导致死锁
内容的提问来源于stack exchange,提问作者feedback
相关产品推荐
相关产品推荐

