Dask多阶段资源配置触发Failed to Serialize序列化错误
问题背景
参照Dask Jobqueue官方示例编写多阶段任务资源调度逻辑,为避免后续文档变更导致代码无法溯源,本次运行的完整代码如下:
from dask_jobqueue import SLURMCluster from distributed import Client from dask import delayed cluster = SLURMCluster(memory='8g', processes=1, cores=2, extra=['--resources ssdGB=200,GPU=2']) cluster.scale(2) client = Client(cluster) def step_1_w_single_GPU(data): return "Step 1 done for: %s" % data def step_2_w_local_IO(data): return "Step 2 done for: %s" % data stage_1 = [delayed(step_1_w_single_GPU)(i) for i in range(10)] stage_2 = [delayed(step_2_w_local_IO)(s2) for s2 in stage_1] result_stage_2 = client.compute(stage_2, resources={tuple(stage_1): {'GPU': 1}, tuple(stage_2): {'ssdGB': 100}})
运行环境版本信息:
- Python 3.8.10
- dask 2022.2.0
- dask-jobqueue 0.7.3
报错信息
运行代码时触发序列化失败,核心报错为TypeError: can not serialize 'Delayed' object,完整错误栈如下:
distributed.protocol.core - CRITICAL - Failed to Serialize Traceback (most recent call last): File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/protocol/core.py", line 76, in dumps frames[0] = msgpack.dumps(msg, default=_encode_default, use_bin_type=True) File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/msgpack/__init__.py", line 38, in packb return Packer(**kwargs).pack(o) File "msgpack/_packer.pyx", line 294, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 300, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 297, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 229, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 291, in msgpack._cmsgpack.Packer._pack TypeError: can not serialize 'Delayed' object distributed.comm.utils - ERROR - can not serialize 'Delayed' object Traceback (most recent call last): File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/comm/utils.py", line 33, in _to_frames return list(protocol.dumps(msg, **kwargs)) File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/protocol/core.py", line 76, in dumps frames[0] = msgpack.dumps(msg, default=_encode_default, use_bin_type=True) File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/msgpack/__init__.py", line 38, in packb return Packer(**kwargs).pack(o) File "msgpack/_packer.pyx", line 294, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 300, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 297, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 229, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 291, in msgpack._cmsgpack.Packer._pack TypeError: can not serialize 'Delayed' object distributed.batched - ERROR - Error in batched write Traceback (most recent call last): File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/batched.py", line 94, in _background_send nbytes = yield self.comm.write( File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/tornado/gen.py", line 762, in run value = future.result() File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/comm/tcp.py", line 250, in write frames = await to_frames( File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/comm/utils.py", line 50, in to_frames return _to_frames() File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/comm/utils.py", line 33, in _to_frames return list(protocol.dumps(msg, **kwargs)) File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/distributed/protocol/core.py", line 76, in dumps frames[0] = msgpack.dumps(msg, default=_encode_default, use_bin_type=True) File "/opt/eagleseven/pyenv/e7cloudv0/lib/python3.8/site-packages/msgpack/__init__.py", line 38, in packb return Packer(**kwargs).pack(o) File "msgpack/_packer.pyx", line 294, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 300, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 297, in msgpack._cmsgpack.Packer.pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 231, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 229, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 264, in msgpack._cmsgpack.Packer._pack File "msgpack/_packer.pyx", line 291, in msgpack._cmsgpack.Packer._pack TypeError: can not serialize 'Delayed' object
问题原因
报错的核心原因是API版本不匹配:
- 当前使用的2022.2.0版本Dask,
client.compute方法的resources参数只接受任务键字符串作为资源映射的键,不支持直接传入Delayed对象。 - 代码中将
stage_1、stage_2两个存满Delayed对象的列表转元组作为字典键,客户端向调度器提交任务时需要序列化这个资源映射字典,而Delayed对象是客户端侧的任务图节点实例,不支持被msgpack序列化,直接触发报错。 - 直接传Delayed对象作为资源键的写法是2022年10月之后版本Dask新增的语法,参考的示例对应更新版本的Dask,和当前环境版本不兼容。
修复方案
两种方案都可以在当前版本正常运行,按需选择即可:
方案1:改用任务键构造资源映射
每个Delayed对象自带.key属性,是调度器可识别的字符串类型任务键,替换原有元组键即可:
from dask_jobqueue import SLURMCluster from distributed import Client from dask import delayed cluster = SLURMCluster(memory='8g', processes=1, cores=2, extra=['--resources ssdGB=200,GPU=2']) cluster.scale(2) client = Client(cluster) def step_1_w_single_GPU(data): return "Step 1 done for: %s" % data def step_2_w_local_IO(data): return "Step 2 done for: %s" % data stage_1 = [delayed(step_1_w_single_GPU)(i) for i in range(10)] stage_2 = [delayed(step_2_w_local_IO)(s2) for s2 in stage_1] # 用每个任务的字符串key构造资源映射 resource_config = {} for task in stage_1: resource_config[task.key] = {'GPU': 1} for task in stage_2: resource_config[task.key] = {'ssdGB': 100} result_stage_2 = client.compute(stage_2, resources=resource_config)
方案2:用annotate上下文在任务定义时绑定资源(推荐)
使用dask.annotate在构造delayed任务时直接标注资源需求,不需要在compute阶段手动维护资源映射,逻辑更清晰,也不会出现键类型错误:
from dask_jobqueue import SLURMCluster from distributed import Client from dask import delayed, annotate cluster = SLURMCluster(memory='8g', processes=1, cores=2, extra=['--resources ssdGB=200,GPU=2']) cluster.scale(2) client = Client(cluster) def step_1_w_single_GPU(data): return "Step 1 done for: %s" % data def step_2_w_local_IO(data): return "Step 2 done for: %s" % data # 第一阶段任务绑定GPU资源 with annotate(resources={'GPU': 1}): stage_1 = [delayed(step_1_w_single_GPU)(i) for i in range(10)] # 第二阶段任务绑定SSD存储资源 with annotate(resources={'ssdGB': 100}): stage_2 = [delayed(step_2_w_local_IO)(s2) for s2 in stage_1] result_stage_2 = client.compute(stage_2)
内容的提问来源于stack exchange,提问作者michaelgbj
相关产品推荐
相关产品推荐

