如何为字典中每个元素创建multiprocessing.Lock并解决竞态问题?
背景与现有代码
我正在开发一个基于多进程的程序,通过文件缓存提升运行速度。该缓存初始为空,会在数据请求时逐步填充,同时存在一组未处理的原始文件作为缓存文件的来源。当前使用的多进程代码如下:
# file_caches is a multiprocessing.Manager.dict() # file_cache_lock is a multiprocessing.Lock() if file_path in file_caches: # We have a cache file_cache_lock.acquire() cached = file_caches[file_path][:] file_cache_lock.release() data1 = cached[0] data2 = cached[1] elif file_path.exists(): data1 = np.load(file_path) data2 = get_data2() if file_cache_lock.acquire(False) and (file_path not in file_caches): # Non-blocking acquire file_caches[file_path] = (data1, data2) file_cache_lock.release() else: # Load original file data1, data2 = read_and_process(original_file_path) # save data file_path.parent.mkdir(parents=True, exist_ok=True) with open(file_path, "wb") as f: np.save(f, data1, allow_pickle=False) if file_cache_lock.acquire(False) and (file_path not in file_caches): # Non-blocking acquire file_caches[file_path] = (data1, data2) file_cache_lock.release()
问题描述
当两个或多个进程先后请求同一文件时,会出现竞态条件问题:例如进程A发现缓存不存在且缓存文件未创建,便开始处理原始文件并创建缓存文件;此时进程B看到文件已创建但尚未写入完成,会进入elif分支读取不完整的数据,导致错误。
尝试在缓存字典的元组中添加一个multiprocessing.Lock()字段,想避免阻塞其他数据的读写同时解决竞态问题,但遇到错误:Lock objects should only be shared between processes through inheritance。请问是否可以动态创建锁并添加到字典中?或者有更优的解决方案?
解决方案
1. 直接动态添加锁到Manager字典不可行
multiprocessing.Manager创建的共享字典无法直接存储普通的multiprocessing.Lock对象,因为这类锁需要通过进程继承或者由Manager主动创建才能跨进程共享,直接动态添加会触发你遇到的错误,这个思路走不通。
2. 原子文件写入解决读写不完整问题
核心问题是文件创建过程不原子,进程B会看到存在但未写完的文件。可以通过临时文件+原子重命名解决:
- 写入缓存时先写临时文件,确认写入完成后,再用
os.replace()原子重命名为目标文件(多数操作系统中,这个重命名操作是原子的)。 - 修改代码中的保存逻辑:
# 替换原save data部分代码 file_path.parent.mkdir(parents=True, exist_ok=True) # 先写入临时文件 temp_path = file_path.with_suffix('.tmp') with open(temp_path, "wb") as f: np.save(f, data1, allow_pickle=False) # 原子重命名,确保其他进程看到的要么是完整文件,要么不存在 os.replace(temp_path, file_path)
这样进程B只有在文件完全写入后才会检测到file_path存在,不会读取到不完整的数据。
3. 基于Manager的Per-File锁(细粒度并发控制)
如果需要对单个文件的操作做更严格的并发控制,可以通过multiprocessing.Manager单独创建一个锁字典,每个文件路径对应一个Manager生成的锁:
# 主进程初始化 manager = multiprocessing.Manager() file_caches = manager.dict() file_cache_lock = manager.Lock() # 全局锁用于保护锁字典的创建 file_locks = manager.dict() # 存储每个文件对应的锁 # 获取文件专属锁的辅助函数 def get_file_lock(file_path): with file_cache_lock: if file_path not in file_locks: file_locks[file_path] = manager.Lock() return file_locks[file_path]
之后处理文件时,先获取对应文件的锁,再执行操作:
file_lock = get_file_lock(file_path) file_lock.acquire() try: if file_path in file_caches: cached = file_caches[file_path] data1, data2 = cached[0], cached[1] elif file_path.exists(): data1 = np.load(file_path) data2 = get_data2() file_caches[file_path] = (data1, data2) else: data1, data2 = read_and_process(original_file_path) # 原子写入缓存文件 file_path.parent.mkdir(parents=True, exist_ok=True) temp_path = file_path.with_suffix('.tmp') with open(temp_path, "wb") as f: np.save(f, data1, allow_pickle=False) os.replace(temp_path, file_path) file_caches[file_path] = (data1, data2) finally: file_lock.release()
这种方式既避免了全局锁的大粒度阻塞,又通过Manager创建的锁实现了跨进程共享,结合原子写入彻底解决竞态问题。
4. 全局锁简化实现(适合低并发场景)
如果你的并发量不大,也可以直接用全局锁包裹整个文件处理逻辑,虽然粒度稍大,但实现简单:
file_cache_lock.acquire() try: if file_path in file_caches: cached = file_caches[file_path] data1, data2 = cached[0], cached[1] elif file_path.exists(): data1 = np.load(file_path) data2 = get_data2() file_caches[file_path] = (data1, data2) else: data1, data2 = read_and_process(original_file_path) # 原子写入缓存文件 file_path.parent.mkdir(parents=True, exist_ok=True) temp_path = file_path.with_suffix('.tmp') with open(temp_path, "wb") as f: np.save(f, data1, allow_pickle=False) os.replace(temp_path, file_path) file_caches[file_path] = (data1, data2) finally: file_cache_lock.release()
这种方式确保同一时间只有一个进程处理某个文件,不会出现竞态,但多个进程处理不同文件时会被全局锁阻塞,适合并发场景不复杂的情况。
内容的提问来源于stack exchange,提问作者Menotdan

