Ray Data 性能调优实战指南:转换、读取、内存与执行配置的全面调优

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 内置的分布式数据处理引擎用于在大规模集群上执行 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_size128) # 不推荐逐行调用 Python 函数无法向量化 # ds ray.data.range(1000).map(lambda x: x * 2)需要注意如果你的转换本身不是向量化的例如依赖逐行逻辑的 Python 函数那么使用map_batches并不会带来性能收益。两者的取舍取决于转换能否以批为单位批量执行。启用 Polars 加速排序use_polars_sortRay 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.pyDEFAULT_USE_POLARS_SORT False因此默认仍走 PyArrow 路径需要手动开启才能切换到 Polars。优化读取输出块数量、资源与列裁剪调优 read 输出块read_output_blocks默认情况下Ray Data 自动为读取操作选择输出块数量具体遵循以下流程传给 Ray Data 读取 API 的override_num_blocks参数指定输出块数量它等价于要创建的读取任务数量。通常如果读取操作后面紧跟map或map_batchesmap 会与读取融合fusion因此override_num_blocks也决定了 map 任务的数量。当未显式指定时Ray Data 按以下启发式规则依次应用决定默认输出块数量以默认值 200 起步。可通过设置DataContext.read_op_min_num_blocks覆盖该常量定义于 python/ray/data/context.pyDEFAULT_READ_OP_MIN_NUM_BLOCKS 200。最小块大小默认 1 MiB。如果块数量会导致块小于该阈值则减少块数量以避免微小块带来的开销。可通过设置DataContext.target_min_block_size字节覆盖默认值见 python/ray/data/context.pyDEFAULT_TARGET_MIN_BLOCK_SIZE 1 * 1024 * 1024。最大块大小默认 128 MiB。如果块数量会导致块大于该阈值则增加块数量以避免处理过程中出现内存不足OOM。可通过设置DataContext.target_max_block_size字节覆盖默认值见 python/ray/data/context.pyDEFAULT_TARGET_MAX_BLOCK_SIZE 128 * 1024 * 1024。可用 CPU。增加块数量以充分利用集群中所有可用 CPU——Ray Data 选择的读取任务数量至少为可用 CPU 数的 2 倍。手动指定override_num_blocks的实战示例有时候手动调优块数量对应用更有利。例如下面的代码把多个文件合并到同一个读取任务中以避免产生过大的块import ray # 假装有两个 CPU。 ray.init(num_cpus2) # 将 iris.csv 重复 16 次。 ds ray.data.read_csv([s3://anonymousray-example-data/iris.csv] * 16) print(ds.materialize())输出示意MaterializedDataset( num_blocks4, num_rows2400, ... )但假设你明确知道希望并行读取全部 16 个文件——例如你预期自动扩缩器autoscaler会向集群添加更多 CPU或者你希望下游算子并行处理每个文件的内容。此时可以通过设置override_num_blocks参数获得该行为。注意下面的代码中输出块数量等于override_num_blocksimport ray # 假装有两个 CPU。 ray.init(num_cpus2) # 将 iris.csv 重复 16 次。 ds ray.data.read_csv([s3://anonymousray-example-data/iris.csv] * 16, override_num_blocks16) print(ds.materialize())输出示意MaterializedDataset( num_blocks16, num_rows2400, ... )自动分块下的块数并非精确保证使用默认的自动检测块数量时Ray Data 试图把每个任务的输出控制在DataContext.target_max_block_size字节以内。但 Ray Data 无法完美预测每个任务输出的大小因此每个任务可能产生一个或多个输出块。这意味着最终Dataset中的总块数可能与指定的override_num_blocks不同。例如手动指定override_num_blocks1但单个任务仍然产出了多个块import ray # 假装有两个 CPU。 ray.init(num_cpus2) # 生成约 400MB 的数据。 ds ray.data.range_tensor(5_000, shape(10_000, ), override_num_blocks1) print(ds.materialize())输出示意MaterializedDataset( num_blocks3, num_rows5000, schema{data: ArrowTensorTypeV2(shape(10000,), dtypeint64)} )输入文件数对读取任务数的上限约束目前 Ray Data 对每个输入文件最多分配一个读取任务。因此如果输入文件数量小于override_num_blocks读取任务数量会被限制为输入文件数。为了保证下游转换仍能以期望的块数执行Ray Data 会将读取任务的输出拆分成总计override_num_blocks个块并阻止与下游转换的融合。换句话说每个读取任务的输出块会先物化到 Ray 对象存储中然后才执行消费它的 map 任务。例如下面的代码只用 1 个任务执行read_csv但在执行map之前其输出被拆分为 4 个块import ray # 假装有两个 CPU。 ray.init(num_cpus2) ds ray.data.read_csv(s3://anonymousray-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_blocksimport ray # 假装有两个 CPU。 ray.init(num_cpus2) ds ray.data.read_csv(s3://anonymousray-example-data/iris.csv, override_num_blocks1).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_cpus0.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://anonymousray-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 参数的整体说明。减少内存使用避免溢写与合并过小块避免对象溢写spillingDataset 的中间块与输出块存放在 Ray 的对象存储中。虽然 Ray Data 通过流式执行streaming execution尽量减少对象存储占用但当工作集超过对象存储容量时Ray 会开始把块溢写spill到磁盘这可能显著拖慢执行速度甚至引发磁盘空间不足错误。有两种场景下溢写是预期行为使用了全对全all-to-allshuffle 操作调用了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_sizebatch_size)并把batch_size设为每个块期望的行数。这种方式以流式执行避免物化。使用map_batches时Ray Data 会合并coalesce块使每个 map 任务至少能处理这么多行。需要注意batch_size是任务输入块大小的下界但并不必然决定任务的最终输出块大小。下面的代码用两种策略把 10 个各含 1 行的小块合并为 1 个含 10 行的大块import ray # 假装有两个 CPU。 ray.init(num_cpus2) # 1. 使用 ds.repartition()。 ds ray.data.range(10, override_num_blocks10).repartition(1) print(ds.materialize().stats()) # 2. 使用 ds.map_batches()。 ds ray.data.range(10, override_num_blocks10).map_batches(lambda batch: batch, batch_size10) 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( cpu10, gpu5, object_store_memory10e9, )从源码看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_cpusParquet 读取显式投影列。内存层用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),仅供参考