如何序列化MetPy(Pint)单位以适配Dask Distributed?
解决MetPy单位在Dask Distributed中的序列化问题
你遇到的TypeError: cannot pickle 'weakref' object错误,本质是Dask序列化任务函数时,尝试打包函数依赖的全局units对象(MetPy的Pint UnitRegistry实例),而这个实例内部包含无法被pickle处理的弱引用结构。
下面是两种可行的解决方案:
方案1:在任务函数内部导入units
把units的导入逻辑移到函数内部,让每个Worker在执行任务时自行初始化units实例,避免序列化全局的units对象:
import metpy.calc as mpcalc from dask.distributed import Client, LocalCluster def calculate_dewpoint(vapor_pressure): # 将units导入放到函数内部 from metpy.units import units dewpoint = mpcalc.dewpoint(vapor_pressure * units('hPa')) return dewpoint cluster = LocalCluster() client = Client(cluster) # 本地运行正常 vapor_pressure = 5 dp = calculate_dewpoint(vapor_pressure) print(dp) # 分布式运行现在可以正常工作 vapor_pressure = 5 dp_future = client.submit(calculate_dewpoint, vapor_pressure) dp = dp_future.result()
方案2:提前在客户端处理单位,传递带单位的数值
如果不想修改函数结构,可以在客户端先给数值加上单位,再传递给任务函数。带单位的Pint Quantity对象本身支持正确序列化,不会触发弱引用错误:
import metpy.calc as mpcalc from metpy.units import units from dask.distributed import Client, LocalCluster def calculate_dewpoint(vapor_pressure): # 直接接收带单位的参数 dewpoint = mpcalc.dewpoint(vapor_pressure) return dewpoint cluster = LocalCluster() client = Client(cluster) # 本地运行正常 vapor_pressure = 5 * units('hPa') dp = calculate_dewpoint(vapor_pressure) print(dp) # 分布式运行正常 vapor_pressure = 5 * units('hPa') dp_future = client.submit(calculate_dewpoint, vapor_pressure) dp = dp_future.result()
内容的提问来源于stack exchange,提问作者bwc
相关产品推荐
相关产品推荐

