如何在含pool.imap_unordered的Python多进程脚本中正确使用Lock()
在multiprocessing.Pool.imap_unordered中正确使用Lock的方法
首先明确:只有当你的my_func中存在多进程共享资源访问(比如写入同一个文件、修改共享内存变量、避免控制台输出混乱)时,才需要使用锁。如果只是单独处理每个DataFrame分组,不需要加锁。
下面是针对你的脚本的两种正确实现方式:
方式一:使用Manager传递共享锁
通过multiprocessing.Manager()创建可跨进程共享的锁,将锁与每个分组数据打包成元组传入任务:
import multiprocessing import pandas as pd def my_func(args): df, lock = args # 临界区操作:需要加锁保护的代码块 with lock: print(f"正在处理分组 a={df['a'].iloc[0]}") # 处理DataFrame逻辑 df = df.head(1) return df if __name__ == "__main__": df = pd.DataFrame({'a': [2, 2, 1, 1, 3, 3], 'b': [4, 5, 6, 4, 5, 6], 'c': [4, 5, 6, 4, 5, 6]}) with multiprocessing.Manager() as manager: lock = manager.Lock() # 将每个分组和锁打包为参数 groups_with_lock = ((g, lock) for _, g in df.groupby("a")) with multiprocessing.Pool() as pool: out = [] for res in pool.imap_unordered(my_func, groups_with_lock): out.append(res) final_df = pd.concat(out) print(final_df)
方式二:通过Pool初始化器传递锁
利用Pool的initializer和initargs参数,在每个子进程启动时注入锁,避免每次任务都传递锁参数:
import multiprocessing import pandas as pd # 全局变量用于存储子进程中的锁引用 lock = None def init_worker(shared_lock): global lock lock = shared_lock def my_func(df): # 使用全局锁保护临界区 with lock: print(f"正在处理分组 a={df['a'].iloc[0]}") # 处理DataFrame逻辑 df = df.head(1) return df if __name__ == "__main__": df = pd.DataFrame({'a': [2, 2, 1, 1, 3, 3], 'b': [4, 5, 6, 4, 5, 6], 'c': [4, 5, 6, 4, 5, 6]}) lock = multiprocessing.Lock() # 创建进程池时传入初始化函数和锁参数 with multiprocessing.Pool(initializer=init_worker, initargs=(lock,)) as pool: groups = (g for _, g in df.groupby("a")) out = [] for res in pool.imap_unordered(my_func, groups): out.append(res) final_df = pd.concat(out) print(final_df)
关键注意点
- 不能直接传递普通的
multiprocessing.Lock()给进程池任务:子进程会复制父进程内存,锁的状态无法跨进程同步,必须用Manager.Lock()或初始化器方式传递共享锁实例。 with lock:上下文管理器会自动处理锁的获取和释放,避免手动操作导致死锁。imap_unordered返回结果的顺序是任务完成顺序,锁的使用不影响结果拼接,最终pd.concat(out)仍能得到正确的合并结果。
内容的提问来源于stack exchange,提问作者Daneel Ank
相关产品推荐
相关产品推荐

