WSL2使用Ray DatasetPipeline与Actor时节点故障致对象丢失问题
问题解答
核心结论
你遇到的问题是Ray 1.x版本的已知限制和内核缺陷共同导致的,你的使用方式不符合对应版本DatasetPipeline的设计规范:
- Ray 1.x版本中DatasetPipeline是绑定创建进程上下文的有状态迭代器,不支持序列化后传递给远程Actor/任务,这是你代码的使用问题
- 跨进程传递Pipeline触发了Ray内核的未处理边界case,进而引发段错误和对象丢失,属于Ray本身的缺陷,该问题在Ray 2.0及以上版本已经修复
根因说明
DatasetPipeline是Ray Data的流式执行抽象,你使用的版本中,Pipeline实例本身持有创建进程内的对象引用、执行调度状态等上下文信息,序列化传递到远端Actor后,这些上下文信息会丢失,导致远端Actor迭代Pipeline时无法找到对应的数据块,触发ObjectLostError。同时零大小数据块的边界 case 没有被Ray C++内核处理,直接触发了段错误,导致Worker进程异常退出。
redis的ulimit警告和Dashboard的aiohttp报错和本次核心问题无关:
- ulimit警告属于环境配置问题,你可以按照提示执行
ulimit -n 65535调整文件句柄限制,避免后续大规模任务出现连接不足问题 - Dashboard报错是日志读取接口的小缺陷,不影响核心任务执行,后续版本已经修复
修复方案
方案1:修改代码适配当前版本
把Pipeline的创建逻辑移到Actor内部执行,避免跨进程传递Pipeline实例,修改后的代码如下:
import ray from typing import List from ray.data.dataset_pipeline import DatasetPipeline data_files = ["data/data1.csv", "data/data2.csv", "data/data3.csv"] @ray.remote class RemotePipelineActor: def __init__(self, shard_idx: int, total_shards: int) -> None: # 在Actor进程内部创建并拆分Pipeline,上下文完全绑定当前进程 full_pipeline = ray.data.read_csv(data_files).pipeline(parallelism=3) self.pipeline = full_pipeline.split(total_shards)[shard_idx] def log_from_pipeline(self) -> int: for df in self.pipeline.iter_batches(batch_size=1000, batch_format="pandas"): pass return 1 def main_pipeline_actor(): total_shards = 3 actors = [RemotePipelineActor.remote(i, total_shards) for i in range(total_shards)] results = ray.get([actor.log_from_pipeline.remote() for actor in actors]) return results if __name__ == "__main__": ray.init() main_pipeline_actor()
方案2:升级Ray版本
升级到Ray 2.0及以上版本,官方已经重构了DatasetPipeline的序列化逻辑,支持跨进程传递,你的原有代码可以直接正常运行。
内容的提问来源于stack exchange,提问作者Anthony Naddeo
相关产品推荐
相关产品推荐

