Ray Data 文本数据处理实战指南:从读取、转换到推理与保存

发布时间:2026/9/19 12:31:33
Ray Data 文本数据处理实战指南:从读取、转换到推理与保存 Ray Data 文本数据处理实战指南从读取、转换到推理与保存【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/rayRay 的 Data 模块ray.data将大规模文本数据处理抽象为分布式数据集Dataset你可以用几行代码完成 TB 级文本的读取、并行转换、批量推理与落盘。本文以仓库文档 working-with-text.rst 为主线结合源码 read_api.py 与 dataset.py 的实现细节系统讲解文本数据全流程读完即可在本地或云存储上搭建一套可复用的文本 ETL 与离线推理流水线。读取文本文件Ray Data 提供三种读取文本的路径按数据形态选用纯文本按行读取、JSONL 按记录读取、任意二进制格式手动解码。逐行读取read_textray.data.read_text将文本文件中的每一行作为一个数据行schema 中默认列名为text类型为stringimport ray ds ray.data.read_text(s3://anonymousray-example-data/this.txt) ds.show(3)输出{text: The Zen of Python, by Tim Peters} {text: Beautiful is better than ugly.} {text: Explicit is better than implicit.}从源码 text_datasource.py 可以看到底层实现TextDatasource._read_stream将整个文件读取后按encoding解码再调用splitlines()逐行切分跳过空行后构造{text: line}记录最后交给DelegatingBlockBuilder组装为数据块Block。这正是一行一记录语义的来源。read_text的完整签名位于 read_api.py常用参数如下参数默认值说明paths必填单个文件、单个目录或包含文件/目录的路径列表encodingutf-8文件编码如utf-8、asciidrop_empty_linesTrue是否丢弃空行设为False时空行也会保留为记录filesystem自动推断PyArrow 文件系统实现不指定时按路径 scheme 自动选择如s3://用S3FileSysteminclude_pathsFalse设为True时把每条记录来源文件路径写入path列ignore_missing_pathsFalse设为True时忽略不存在的路径shuffleNone设为files或FileShuffleConfig(seed...)可随机打乱输入文件顺序file_extensionsNone按扩展名过滤待读文件num_cpus/num_gpusNone每个并行读取 worker 预留的 CPU/GPU 数量memoryNone每个并行读取 worker 预留的堆内存字节concurrency动态决定并发执行的 Ray 任务上限override_num_blocks动态决定覆盖输出数据块数量绝大多数场景无需手动设置label_selector/fallback_strategyNone节点标签约束及候选回退策略参数行为在仓库测试 test_text.py 中得到了完整验证test_read_text验证多文件拼接与drop_empty_lines语义test_read_text_remote_args验证通过ray_remote_args{resources: {...}}把读取任务调度到特定节点test_fsspec_http_file_system验证不指定filesystem时 HTTP 文件系统的自动解析。按记录读取JSON LinesJSON Lines 是每行一个 JSON 对象的文本格式适合逐条流式处理结构化记录。调用ray.data.read_json时每个 JSON 对象成为一行记录import ray ds ray.data.read_json(s3://anonymousray-example-data/logs.json) ds.show(3)输出{timestamp: datetime.datetime(2022, 2, 8, 15, 43, 41), size: 48261360} {timestamp: datetime.datetime(2011, 12, 29, 0, 19, 10), size: 519523} {timestamp: datetime.datetime(2028, 9, 9, 5, 6, 7), size: 2163626}注意 read_api.py 中lines参数默认False对于普通 JSON 文件整个文件被读为一行对于 JSONL 文件需传linesTrue才能逐行解析。partitioning默认为Partitioning(hive)即自动从路径解析 Hive 风格分区如year2022/month09/...这可以从 read_json 的 docstring 示例 中得到印证。实现层面json_datasource.py 定义了默认支持的扩展名集合json、jsonl以及 gzip.gz、Brotli.br、Zstandard.zst、lz4.lz4压缩变体读取时按扩展名自动解压。解析优先使用 PyArrow 的 JSON reader遇到跨块边界等ArrowInvalid错误时会自动几何级增大块大小重试若 PyArrow 仍失败则回退到原生json.load()见 json_datasource.py 与 L146-L161这是 JSON 读取健壮性的关键保障。其他格式读取原始二进制再手动解码对于 HTML、XML、PDF 等无内置 reader 的文本格式先用ray.data.read_binary_files以原始字节读入schema 为bytes列再通过Dataset.map手动解码。以 HTML 为例配合 BeautifulSoup 抽取正文from typing import Any, Dict from bs4 import BeautifulSoup import ray def parse_html(row: Dict[str, Any]) - Dict[str, Any]: html row[bytes].decode(utf-8) soup BeautifulSoup(html, featureshtml.parser) return {text: soup.get_text().strip()} ds ( ray.data.read_binary_files(s3://anonymousray-example-data/index.html) .map(parse_html) ) ds.show()输出{text: Batoidea\nBatoidea is a superorder of cartilaginous fishes...}read_binary_files的签名在 read_api.py同样支持include_paths、filesystem、shuffle、num_cpus/num_gpus等参数。当include_pathsTrue时文件路径会以path列形式出现便于后续按来源回溯。更完整的读取方式清单CSV、Parquet、TFRecords、Zarr 等参见 Loading data 指南其中也给出了 Text 与 Binary 读取的 schema 示例。读取本地共享存储如 NFS或云存储S3/GCS 等时务必保证路径在集群所有节点上可见且所有节点已完成云厂商认证。转换文本文本转换只需把处理逻辑封装成普通函数或可调用类然后交给map或map_batchesRay Data 会在分布式任务中并行执行。行级转换map逐行处理用Dataset.map函数接收一个 dict 行、返回一个 dict 行。例如统一小写化from typing import Any, Dict import ray def to_lower(row: Dict[str, Any]) - Dict[str, Any]: row[text] row[text].lower() return row ds ( ray.data.read_text(s3://anonymousray-example-data/this.txt) .map(to_lower) ) ds.show(3)输出{text: the zen of python, by tim peters} {text: beautiful is better than ugly.} {text: explicit is better than implicit.}批量转换map_batches如果你的转换是向量化的NumPy、pandas 风格map_batches性能更优——它把多行打包成一个批次传入函数一次调用处理整批数据见 map 的 docstring 提示。map_batches的关键参数batch_size支持三种取值源码见 dataset.pyauto由 Ray Data 依据数据与资源自动决定批次行数整数指定每批行数向量化转换下增大批次通常提升吞吐但过大会引发内存/显存溢出OOM遇到 OOM 应调小None把每个数据块整体作为一批。更详细的批次格式batch_formatnumpy/pandas/pyarrow与性能取舍参见 Transforming data 指南。可调用类与状态化转换当转换需要昂贵的初始化如下载模型权重、建立连接池时改用可调用类。__init__在每个 worker 上只执行一次完成装载__call__处理数据函数则是无状态的任何初始化都会对每条数据重复执行。内部实现上函数走 Ray Task类走 Ray Actor详见 transforming-data.rst 的 Stateful Transforms 章节。对文本执行推理离线推理是文本处理最常见的下游场景。标准做法是写一个可调用类在__init__中装载预训练模型在__call__中按批次调用模型随后用map_batches传入该类并通过ActorPoolStrategy控制并行度。CPU 文本分类示例from typing import Dict import numpy as np from transformers import pipeline import ray class TextClassifier: def __init__(self): self.model pipeline(text-classification) def __call__(self, batch: Dict[str, np.ndarray]) - Dict[str, list]: predictions self.model(list(batch[text])) batch[label] [prediction[label] for prediction in predictions] return batch ds ( ray.data.read_text(s3://anonymousray-example-data/this.txt) .map_batches(TextClassifier, computeray.data.ActorPoolStrategy(size2), batch_sizeauto) ) ds.show(3)输出{text: The Zen of Python, by Tim Peters, label: POSITIVE} {text: Beautiful is better than ugly., label: POSITIVE} {text: Explicit is better than implicit., label: POSITIVE}compute 与 ActorPoolStrategy 的取值compute参数控制 map/map_batches 的调度方式dataset.py 有完整说明对函数不指定compute使用TaskPoolStrategy()按可用资源与输入块数并发调度 Ray 任务对可调用类不指定compute使用ActorPoolStrategy(min_size1, max_sizeNone)启动从 1 到不限的自伸缩 Actor 池ActorPoolStrategy(sizen)固定 n 个 ActorActorPoolStrategy(min_sizem, max_sizen, initial_sizeinitial)从 m 到 n 自伸缩、初始为 initial 的 Actor 池。由于每个 Actor 都持有模型副本size应结合集群资源设定。对于加载大模型这类装载昂贵、调用便宜的场景Actor 复用避免了反复初始化这正是离线推理的推荐形态。上 GPUnum_gpus 与显式 batch_size若使用 GPU 推理需要三处改动详见 batch_inference.rst 的 Using GPUs 章节类实现中把模型与输入显式搬到 GPU如pipeline(..., devicecuda:0)或model.cuda()在map_batches中传num_gpus1表示每个 Actor 占 1 张 GPU指定整数batch_sizeGPU 场景下batch_sizeauto不被允许见 dataset.py 的校验逻辑。CPU 推理推荐batch_sizeautoGPU 推理则在显存允许范围内取尽可能大的整数批次以提升 GPU 利用率发生 OOM 时下调即可。并发度通常与集群 GPU 数对齐例如ActorPoolStrategy(size2)num_gpus1共占用 2 张 GPU。完整端到端流程作业级 checkpoint、OOM 排查等见 End-to-end: Offline Batch Inference。面向大语言模型LLM的规模化批量推理vLLM/SGLang 引擎、多 GPU 分片、OpenAI 兼容端点等已独立成体系参见 Working with LLMs。保存文本处理后的结果可用各类write_*方法落盘Ray Data 支持 Parquet、JSON、CSV 等多种格式。以最常见的 Parquet 为例import ray ds ray.data.read_text(s3://anonymousray-example-data/this.txt) ds.write_parquet(s3://my-bucket/results)write_parquetdataset.py的关键参数包括path目标根目录写出的文件数与数据集块数一致可用repartition调整partition_cols按列做 Hive 风格分区写出mode写入模式默认SaveMode.APPEND可显式指定overwrite等filesystem默认按路径 scheme 自动选择如s3://→S3FileSystem也可显式传入如 GCS 用gcsfs的GCSFileSystem配合gcs://filename_provider自定义输出文件名默认形如{uuid}_{block_idx}.parquetmin_rows_per_file/max_rows_per_file控制单文件行数arrow_parquet_args/arrow_parquet_args_fn透传给 PyArrowParquetWriter的参数。写目标的选型要点见 Saving data 指南共享本地存储使用 NFS 等共享文件系统并挂载到所有 Ray 节点同一路径然后指定该挂载目录云存储先让所有节点完成云厂商认证再用带 scheme 的 URI 指向桶或目录s3://my-bucket/my-folder注意不要使用已废弃的local://scheme。若后续仍需继续处理Parquet 这类列式格式在读取时可利用列裁剪等特性提升效率更多读写 API 详见 Loading data 与 Saving data。小结围绕ray.data文本处理流水线可归纳为一条主线读read_text按行→read_jsonJSONL 记录→read_binary_filesmap任意格式手动解码转无状态逻辑用函数 map/map_batches昂贵初始化用可调用类 ActorPoolStrategy推类内装载模型map_batches配合num_gpus/batch_size做 CPU 或 GPU 离线推理规模化的 LLM 推理可转向 ray.data.llm 体系存write_parquet等写方法落到本地共享存储或云存储。每个环节的底层行为都有明确的源码与测试支撑文本按行切分的细节见 text_datasource.pyJSON 解析的回退机制见 json_datasource.py读写 API 的全量参数见 read_api.py 与 dataset.py行为断言见 test_text.py。按此流程即可在 Ray 集群上构建从原始文本到模型结果的全分布式处理链路。【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考