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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 20:00:28