修改Python模块变量后为何无法传递至新并行进程?
问题原因
当你用Dask分布式执行任务时,worker是独立的Python进程,和主进程完全隔离。你传递的conf模块对象在序列化(通过cloudpickle)并发送给worker时,并不会把你修改后的属性值一起传递——Dask序列化模块对象的逻辑是在worker进程中重新导入该模块,而不是复制主进程中已经修改过的模块实例。所以worker拿到的是configuration.py文件里原始的result_folder = "aFolder",自然输出的是原始值。
解决方案
不要直接传递模块作为配置载体,改用字典或者数据类这种可以完整序列化并传递修改后状态的结构。
方法1:使用字典传递配置
修改你的代码如下:
串行模式
def embarassing(x, conf): print(x) print(conf["result_folder"]) # ... 执行复杂操作并返回值 if __name__ == "__main__": # 用字典代替模块存储配置 conf = {"result_folder": "aFolder"} x = 1 y = embarassing(x, conf) # 临时任务修改配置 x = 2 conf["result_folder"] = "newFold" y = embarassing(x, conf)
并行模式
from dask.distributed import Client def embarassing(x, conf): print(x) print(conf["result_folder"]) # ... 执行复杂操作并返回值 if __name__ == "__main__": conf = {"result_folder": "aFolder"} client = Client(n_workers=1) # 原始任务 x = 1 future = client.submit(embarassing, x, conf) y = future.result() # 临时任务修改配置 x = 2 conf["result_folder"] = "newFold" future = client.submit(embarassing, x, conf) y = future.result()
这样修改后,字典会被完整序列化并传递给worker,修改后的result_folder值能正确生效,并行模式下会输出2和newFold。
方法2:使用数据类(更结构化)
如果你需要更结构化的配置,可以用Python的dataclasses:
from dataclasses import dataclass from dask.distributed import Client @dataclass class Config: result_folder: str def embarassing(x, conf: Config): print(x) print(conf.result_folder) # ... 执行复杂操作并返回值 if __name__ == "__main__": conf = Config(result_folder="aFolder") client = Client(n_workers=1) # 临时任务修改配置 x = 2 conf.result_folder = "newFold" future = client.submit(embarassing, x, conf) y = future.result()
数据类同样可以被Dask正确序列化,传递修改后的属性值到worker进程。
内容的提问来源于stack exchange,提问作者OneTimeShot
相关产品推荐
相关产品推荐

