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/ray
Ray Data 是 Ray 内置的分布式数据处理引擎,用于在大规模集群上执行 ETL、数据预处理与训练数据管道。本文以 Ray Data 的性能调优主题为主线,系统讲解如何优化map类转换、开启 Polars 排序、调优读取任务的输出块数量与资源占用、通过 Parquet 列裁剪(projection pushdown)减少 IO,以及如何控制对象存储溢写(spilling)、合并过小数据块,并通过全局DataContext配置执行资源与确定性执行。读完本文,你将掌握一套可直接落地的 Ray Data 调优工具箱,并理解每个参数在 python/ray/data/context.py 等源码中的真实实现与默认值,从而针对自己的数据集做出有依据的调优决策。
优化转换:批处理优先,必要时启用 Polars
使用map_batches而非map处理向量化转换
如果你的转换是向量化的——例如大多数 NumPy 或 pandas 运算——请使用ray.data.Dataset.map_batches而不是ray.data.Dataset.map。前者的输入是整批数据(batch),允许底层向量化库一次处理多行,开销更低,因此更快。
import ray # 推荐:向量化转换按批次处理 ds = ray.data.range(1000).map_batches(lambda batch: batch * 2, batch_size=128) # 不推荐:逐行调用 Python 函数,无法向量化 # ds = ray.data.range(1000).map(lambda x: x * 2)需要注意:如果你的转换本身不是向量化的(例如依赖逐行逻辑的 Python 函数),那么使用map_batches并不会带来性能收益。两者的取舍取决于转换能否以批为单位批量执行。
启用 Polars 加速排序:use_polars_sort
Ray Data 的sort以及内部需要排序的操作(例如GroupedData.map_groups)默认使用 PyArrow 完成排序步骤。对于大型表格数据集,你可以通过开启 Polars 来加速内部排序:
import ray ctx = ray.data.DataContext.get_current() ctx.use_polars_sort = True开启该标志后,Ray Data 在内部排序步骤中使用 Polars 替代 PyArrow;该标志不影响map_batches等其他操作。
从源码层面看,这个开关在 python/ray/data/_internal/arrow_block.py 中生效:get_sort_transform(context)与get_concat_and_sort_transform(context)会根据context.use_polars or context.use_polars_sort选择 transform_polars.py 中的sort/concat_and_sort实现,否则回退到 transform_pyarrow.py。该标志的默认值定义在 python/ray/data/context.py(DEFAULT_USE_POLARS_SORT = False),因此默认仍走 PyArrow 路径,需要手动开启才能切换到 Polars。
优化读取:输出块数量、资源与列裁剪
调优 read 输出块(read_output_blocks)
默认情况下,Ray Data 自动为读取操作选择输出块数量,具体遵循以下流程:
- 传给 Ray Data 读取 API 的
override_num_blocks参数指定输出块数量,它等价于要创建的读取任务数量。 - 通常,如果读取操作后面紧跟
map或map_batches,map 会与读取融合(fusion),因此override_num_blocks也决定了 map 任务的数量。
当未显式指定时,Ray Data 按以下启发式规则(依次应用)决定默认输出块数量:
- 以默认值 200 起步。可通过设置
DataContext.read_op_min_num_blocks覆盖,该常量定义于 python/ray/data/context.py(DEFAULT_READ_OP_MIN_NUM_BLOCKS = 200)。 - 最小块大小(默认 1 MiB)。如果块数量会导致块小于该阈值,则减少块数量,以避免微小块带来的开销。可通过设置
DataContext.target_min_block_size(字节)覆盖,默认值见 python/ray/data/context.py(DEFAULT_TARGET_MIN_BLOCK_SIZE = 1 * 1024 * 1024)。 - 最大块大小(默认 128 MiB)。如果块数量会导致块大于该阈值,则增加块数量,以避免处理过程中出现内存不足(OOM)。可通过设置
DataContext.target_max_block_size(字节)覆盖,默认值见 python/ray/data/context.py(DEFAULT_TARGET_MAX_BLOCK_SIZE = 128 * 1024 * 1024)。 - 可用 CPU。增加块数量以充分利用集群中所有可用 CPU——Ray Data 选择的读取任务数量至少为可用 CPU 数的 2 倍。
手动指定override_num_blocks的实战示例
有时候手动调优块数量对应用更有利。例如,下面的代码把多个文件合并到同一个读取任务中,以避免产生过大的块:
import ray # 假装有两个 CPU。 ray.init(num_cpus=2) # 将 iris.csv 重复 16 次。 ds = ray.data.read_csv(["s3://anonymous@ray-example-data/iris.csv"] * 16) print(ds.materialize())输出(示意):
MaterializedDataset( num_blocks=4, num_rows=2400, ... )但假设你明确知道希望并行读取全部 16 个文件——例如你预期自动扩缩器(autoscaler)会向集群添加更多 CPU,或者你希望下游算子并行处理每个文件的内容。此时可以通过设置override_num_blocks参数获得该行为。注意下面的代码中,输出块数量等于override_num_blocks:
import ray # 假装有两个 CPU。 ray.init(num_cpus=2) # 将 iris.csv 重复 16 次。 ds = ray.data.read_csv(["s3://anonymous@ray-example-data/iris.csv"] * 16, override_num_blocks=16) print(ds.materialize())输出(示意):
MaterializedDataset( num_blocks=16, num_rows=2400, ... )自动分块下的块数并非精确保证
使用默认的自动检测块数量时,Ray Data 试图把每个任务的输出控制在DataContext.target_max_block_size字节以内。但 Ray Data 无法完美预测每个任务输出的大小,因此每个任务可能产生一个或多个输出块。这意味着最终Dataset中的总块数可能与指定的override_num_blocks不同。例如,手动指定override_num_blocks=1,但单个任务仍然产出了多个块:
import ray # 假装有两个 CPU。 ray.init(num_cpus=2) # 生成约 400MB 的数据。 ds = ray.data.range_tensor(5_000, shape=(10_000, ), override_num_blocks=1) print(ds.materialize())输出(示意):
MaterializedDataset( num_blocks=3, num_rows=5000, schema={data: ArrowTensorTypeV2(shape=(10000,), dtype=int64)} )输入文件数对读取任务数的上限约束
目前 Ray Data 对每个输入文件最多分配一个读取任务。因此,如果输入文件数量小于override_num_blocks,读取任务数量会被限制为输入文件数。为了保证下游转换仍能以期望的块数执行,Ray Data 会将读取任务的输出拆分成总计override_num_blocks个块,并阻止与下游转换的融合。换句话说,每个读取任务的输出块会先物化到 Ray 对象存储中,然后才执行消费它的 map 任务。
例如,下面的代码只用 1 个任务执行read_csv,但在执行map之前其输出被拆分为 4 个块:
import ray # 假装有两个 CPU。 ray.init(num_cpus=2) ds = ray.data.read_csv("s3://anonymous@ray-example-data/iris.csv").map(lambda row: row) print(ds.materialize().stats())输出(示意):
... Operator 1 ReadCSV->SplitBlocks(4): 1 tasks executed, 4 blocks produced in 0.01s ... Operator 2 Map(<lambda>): 4 tasks executed, 4 blocks produced in 0.3s ...要关闭这种行为并允许读取与 map 算子融合,请手动设置override_num_blocks。例如下面的代码让文件数等于override_num_blocks:
import ray # 假装有两个 CPU。 ray.init(num_cpus=2) ds = ray.data.read_csv("s3://anonymous@ray-example-data/iris.csv", override_num_blocks=1).map(lambda row: row) print(ds.materialize().stats())输出(示意):
... Operator 1 ReadCSV->Map(<lambda>): 1 tasks executed, 1 blocks produced in 0.01s ...可以看到,此时读取与 map 被融合为单个算子,避免了中间块的物化开销。
调优读取资源(tuning_read_resources)
默认情况下,Ray 为每个读取任务请求 1 个 CPU,这意味着每个 CPU 同一时刻只能并发执行一个读取任务。对于受益于更高 IO 并行的数据源,可以为每个读取任务保留更少的 CPU。例如,使用ray.data.read_parquet(path, num_cpus=0.25)可以让每个 CPU 上并发执行最多 4 个读取任务,从而在 IO 密集场景下提升吞吐。
Parquet 列裁剪(projection pushdown)
默认情况下,ray.data.read_parquet会把 Parquet 文件中的所有列都读入内存。如果你只需要其中一部分列,请在调用read_parquet时显式指定列清单,以避免加载不必要的数据(即投影下推 / projection pushdown)。这比先读入全部列再调用Dataset.select_columns更高效,因为列选择被下推到了文件扫描阶段。原文档给出了一个展示 schema 的示例:
import ray # 读取 Iris 数据集五列中的两列。 ds = ray.data.read_parquet( "s3://anonymous@ray-example-data/iris.parquet", ).select_columns(["sepal.length", "variety"]) print(ds.schema())输出:
Column Type ------ ---- sepal.length double variety string更贴近"下推"语义的写法是直接在读取阶段传入columns参数,例如ray.data.read_parquet(path, columns=["sepal.length", "variety"]),让列裁剪在文件扫描时完成,进一步减少反序列化与内存占用。该优化对以列为存储单位的 Parquet 格式收益尤其明显,也适用于其他支持列裁剪的数据源(可参考 doc/source/data/loading-data.rst 中关于读取 API 参数的整体说明)。
减少内存使用:避免溢写与合并过小块
避免对象溢写(spilling)
Dataset 的中间块与输出块存放在 Ray 的对象存储中。虽然 Ray Data 通过流式执行(streaming execution)尽量减少对象存储占用,但当工作集超过对象存储容量时,Ray 会开始把块溢写(spill)到磁盘,这可能显著拖慢执行速度,甚至引发磁盘空间不足错误。
有两种场景下溢写是预期行为:
- 使用了全对全(all-to-all)shuffle 操作;
- 调用了
ds.materialize()。
除此之外,最好调优应用以避免溢写。推荐的策略是手动增加读取输出块数量(见上文 调优 read 输出块),或修改应用代码确保每个任务读取的数据量更小。
说明:这是 Ray Data 正在积极发展的领域。如果你的 Dataset 发生溢写且原因不明,可以在 Ray Data 的 issue 系统中按
[data]标签提交反馈。
处理过小的块(too-small blocks)
当 Dataset 中不同算子产出的块大小差异很大时,可能会出现非常小的块,这会损害性能,甚至因元数据过多而导致崩溃。使用ds.stats()检查每个算子的输出块是否至少为 1 MB,理想情况下大于 100 MB。
如果块过小,可以考虑重新分区以得到更大的块,有两种方式:
- 精确控制输出块数:使用
ds.repartition(num_partitions)。注意这是全对全(all-to-all)操作,会在执行重分区前把全部块物化到内存中。 - 无需精确控制块数、只想要更大的块:使用
ds.map_batches(lambda batch: batch, batch_size=batch_size),并把batch_size设为每个块期望的行数。这种方式以流式执行,避免物化。
使用map_batches时,Ray Data 会合并(coalesce)块,使每个 map 任务至少能处理这么多行。需要注意,batch_size是任务输入块大小的下界,但并不必然决定任务的最终输出块大小。
下面的代码用两种策略,把 10 个各含 1 行的小块合并为 1 个含 10 行的大块:
import ray # 假装有两个 CPU。 ray.init(num_cpus=2) # 1. 使用 ds.repartition()。 ds = ray.data.range(10, override_num_blocks=10).repartition(1) print(ds.materialize().stats()) # 2. 使用 ds.map_batches()。 ds = ray.data.range(10, override_num_blocks=10).map_batches(lambda batch: batch, batch_size=10) print(ds.materialize().stats())输出(示意):
# 1. ds.repartition() 的输出。 Operator 1 ReadRange: 10 tasks executed, 10 blocks produced in 0.33s ... * Output num rows: 1 min, 1 max, 1 mean, 10 total ... Operator 2 Repartition: executed in 0.36s Suboperator 0 RepartitionSplit: 10 tasks executed, 10 blocks produced ... Suboperator 1 RepartitionReduce: 1 tasks executed, 1 blocks produced ... * Output num rows: 10 min, 10 max, 10 mean, 10 total ... # 2. ds.map_batches() 的输出。 Operator 1 ReadRange->MapBatches(<lambda>): 1 tasks executed, 1 blocks produced in 0s ... * Output num rows: 10 min, 10 max, 10 mean, 10 total从输出可以看到:repartition路径包含RepartitionSplit/RepartitionReduce两个子算子(物化式全对全),而map_batches路径直接把读取与转换融合为单个算子,以流式方式完成合并。
配置执行:资源限制与局部性
默认情况下,CPU 与 GPU 上限设置为集群规模,对象存储内存上限则保守地设置为对象存储总大小的 1/4,以避免磁盘溢写的可能。以下场景可能需要自定义这些限制:
- 在集群上同时运行多个任务时,设置更低的限制可以避免任务之间的资源争抢;
- 想精细调优内存上限以最大化性能时;
- 为训练任务加载数据时,可以把对象存储内存设为较低值(例如 2 GB),以限制资源占用。
可以通过全局DataContext配置执行选项。这些选项会应用于该进程后续启动的任务:
ctx = ray.data.DataContext.get_current() ctx.execution_options.resource_limits = ctx.execution_options.resource_limits.copy( cpu=10, gpu=5, object_store_memory=10e9, )从源码看,ExecutionOptions定义在 python/ray/data/_internal/execution/interfaces/execution_options.py,其resource_limits字段类型为ExecutionResources,默认使用ExecutionResources.for_limits()(即不设限),preserve_order默认值为False。DataContext通过execution_options字段(python/ray/data/context.py)持有这些配置。注意:DataContext的更改应在创建Dataset之前完成,创建后再修改不会生效;配置对象会自动传播到各个 worker,在 driver 与远端 worker 中均可通过DataContext.get_current()访问。
可重现性:确定性执行
默认情况下preserve_order为False。要启用确定性执行,请将其设置为True:
# 默认情况下,该值为 False。 ctx.execution_options.preserve_order = True该设置可能降低性能,但可以保证块在处理过程中保持顺序。该标志默认关闭。如果你的管道对输出顺序敏感(例如下游需要按块序消费),可以按需开启;若更看重吞吐,则保持默认即可。
小结与调优路线
综合全文,Ray Data 性能调优可以归纳为四条主线,按投入产出比排序:
- 转换层:优先使用向量化的
map_batches;大型表格排序开启use_polars_sort。 - 读取层:用
override_num_blocks匹配文件数与期望并行度;IO 密集场景调低num_cpus;Parquet 读取显式投影列。 - 内存层:用
ds.stats()观察块大小,块过小用repartition或map_batches合并;出现非预期溢写时增大读取并行度、减小单任务数据量。 - 执行层:用
DataContext.execution_options设置 CPU / GPU / 对象存储内存上限以适配多任务共享与训练加载场景;需要确定性输出时开启preserve_order。
每一处配置的默认值都可以在 python/ray/data/context.py 中查到,并结合ds.stats()的实际输出反复迭代,最终形成适合你数据规模与集群拓扑的调优方案。更完整的 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/ray
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考