You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.01 12:57:01