Ray Data Shuffle V2 提升性能与可扩展性
Ray Data 推出了 Shuffle V2,这是一个重新设计的 shuffle 引擎,旨在解决分布式数据处理任务(如连接、分组和重新分区)中的可扩展性和性能瓶颈。据 官方公告,该更新通过利用 Ray 对象存储管理 shuffle 中间数据,而非依赖内存密集型的 actor 堆,从而实现了高达 53 倍的操作加速。
Shuffle 过程对于诸如去重和基于键的连接等操作至关重要,但历史上一直是分布式系统中的痛点。较早的 Shuffle V1 实现依赖于长期运行的聚合器 actor,这引入了内存上限、恢复限制和资源低效等问题。Shuffle V2 通过重新设计引擎,将操作分为 map 和 reduce 阶段,并将中间数据存储在支持溢出的、具有谱系追踪的对象存储分片中,从而消除了这些限制。这种架构不仅增强了可靠性,还解锁了诸如矢量化聚合和分片压缩等优化功能。
为什么重要
分布式工作负载在扩展到更大数据集时通常会因内存不足(OOM)错误或资源分配低效而失败。例如,Ray Data 之前的 shuffle 模型需要足够的内存来在聚合器堆中保存整个数据集,对于超过 1 TB 的数据集,这会频繁导致崩溃。通过允许中间数据溢出到磁盘并利用 Ray Core 的谱系重建功能,Shuffle V2 即使在极端数据负载下也能确保稳定性和可恢复性。
性能提升十分显著。在使用规模化的 TPC-H 数据集进行的内部基准测试中,Shuffle V2 在大约 320 秒内完成了关键连接操作,而 Shuffle V1 要么超时,要么需要显著更多的资源。对于分组操作,矢量化聚合通过将计算卸载到 Arrow 的原生分组功能,实现了超过 30 倍的性能提升。
关键功能和优化
- 可溢出的状态:Shuffle 中间数据被存储为可以溢出到磁盘的 Ray 对象,从而消除了内存限制。
- 压缩:中间数据使用 Zstd 等编解码器进行压缩,减少了存储和传输开销。例如,在测试中,Zstd 将溢出字节减少了 91%,运行时间缩短了 26%。
- 动态资源分配:与 V1 的固定聚合器池不同,V2 根据实际任务需求动态分配资源,从而提高集群效率。
- 融合与矢量化:分组和 map_batches 等下游操作被融合,跳过了不必要的中间写入。聚合现在以列式操作的形式运行,大幅减少了处理时间。
未来发展
Shuffle V2 从 Ray 2.58 开始可用,并计划在 Ray 2.59 中进一步增强。即将推出的功能包括磁盘 shuffle,它完全绕过对象存储以实现更大的可扩展性,以及增量连接支持,这将减少大规模连接的内存开销。此外,当前使用旧引擎的 sort 和 random_shuffle 函数将在未来更新中迁移到 V2 架构。
对用户的意义
对于处理大规模数据管道的团队——尤其是涉及连接、分组或重新分区的任务——Shuffle V2 提供了显著的性能提升和更高的可靠性。用户可以通过以下配置启用该功能:
from ray.data.context import DataContext, ShuffleStrategy
DataContext.get_current().shuffle_strategy = ShuffleStrategy.HASH_SHUFFLE_V2
随着分布式数据工作负载的规模和复杂性不断增加,像 Ray Data Shuffle V2 这样的工具正在为可扩展性和容错处理设定新的基准。凭借高达 53 倍的性能提升以及一个承诺更高可扩展性的路线图,采用 Shuffle V2 对于数据团队来说可能是一个改变游戏规则的选择。