如何在joblib中使用multiprocessing锁?解决锁不可pickle问题
解决joblib中multiprocessing/loky后端无法传递锁的问题
直接传递multiprocessing.Lock给joblib任务会报错,核心原因是原生锁对象无法被pickle序列化,而joblib在分发任务时需要把参数序列化后传给子进程。下面给两种可行的解决办法:
方法一:用multiprocessing.Manager创建可共享锁
multiprocessing.Manager会启动一个独立的管理进程,它创建的锁是代理对象,支持跨进程序列化和访问,完美适配joblib的需求。修改后的代码如下:
from multiprocessing import Process, Lock, Manager from joblib import Parallel, delayed def f(l, i): l.acquire() try: print('hello world', i) finally: l.release() if __name__ == '__main__': # 原生锁用于直接启动Process的场景 lock = Lock() for num in range(10): Process(target=f, args=(lock, num)).start() # 用Manager创建可共享锁给joblib用 with Manager() as manager: shared_lock = manager.Lock() Parallel(n_jobs=2)(delayed(f)(shared_lock, i) for i in range(10, 20))
方法二:针对loky后端使用专属锁(可选)
如果你明确用loky作为joblib后端(joblib默认后端就是loky),可以直接用loky提供的可复用锁,不需要显式传递,子进程能直接获取同一个锁实例:
from multiprocessing import Process, Lock from joblib import Parallel, delayed from loky import get_reusable_lock def f(i): lock = get_reusable_lock("my_shared_lock") lock.acquire() try: print('hello world', i) finally: lock.release() if __name__ == '__main__': lock = Lock() for num in range(10): Process(target=f, args=(num,)).start() Parallel(n_jobs=2, backend="loky")(delayed(f)(i) for i in range(10, 20))
为什么全局变量方案不行?
你之前尝试把锁设为全局变量没用,是因为joblib启动的每个子进程都会重新加载模块,全局变量在每个子进程里都是新的锁实例,不是同一个对象,根本起不到同步作用。
内容的提问来源于stack exchange,提问作者Jann Poppinga
相关产品推荐
相关产品推荐

