2024 年,NVIDIA 与 Polars 团队共同推出 Polars GPU 引擎。它最初基于内存执行器,在单个 GPU 上运行查询。此后,我们反复听到用户希望拥有更大的处理余量:处理更大的数据集,并能够使用多个 GPU。首个实验版本于 25.06 发布,采用由 Dask 编排的流式执行器。
从 cudf-polars 26.06 版本开始,Polars GPU 引擎使用基于 RapidsMPF 的新流式后端。RapidsMPF 是一个用于流式执行和多 GPU 执行的库。新后端重写了早期执行器,速度显著提升,目前是任何非简单 GPU 工作负载的推荐路径。
流式执行器把数据拆成分块,让分块流经查询图,底层数据移动由 RapidsMPF 管理。由此得到两项能力:
- 查询不再受 GPU 显存容量限制。设备内存紧张时,分块会溢写到主机内存,使单个 GPU 可以处理规模为自身容量数倍的数据集。
- 同一个查询可以从单 GPU 扩展到多 GPU。多 GPU 执行只是让同一个执行器在更多 GPU 上运行,无须单独的代码路径。
CPU 引擎针对单节点上的交互式和中等规模分析进行了高度优化,对于大量任务,它仍是合适的工具。但是,当数据集增长到数百 GB 乃至更大时,查询中开销较高的部分——连接、高基数分组聚合和排序——可能显著变慢,而这些正是 GPU 引擎加速最明显的操作。
下文介绍如何运行新引擎、后端如何工作,以及如何用它将 TB 级基准测试最多加速到 23 倍。
安装与配置
引擎通过 cudf-polars 软件包提供。请遵循安装指南,选择与 CUDA 和 Python 版本匹配的构建版本:
# pip, CUDA 12 by default
pip install "polars[gpu]"
# for CUDA 13
pip install cudf-polars-cu13
# or conda
conda install -c rapidsai -c conda-forge cudf-polars
在 .collect() 调用中传入 engine= 参数,即可选择 GPU 执行。最简单的写法是将 engine 设置为 "gpu";这样查询在单 GPU 上执行,无须额外设置:
import polars as pl
query = (
pl.scan_parquet("/data/dataset/*.parquet")
.group_by("customer_id")
.agg(pl.col("amount").sum())
)
result = query.collect(engine="gpu")
要使用多个 GPU 或调整引擎配置,需要构造引擎对象。本文示例使用 RayEngine,也同样支持 DaskEngine 和 SPMDEngine。引擎文档介绍了各自适用的情况。每种引擎都有对应的 pip 可选依赖:
pip install "cudf-polars-cu13[ray]" # or dask
不传入参数构造 RayEngine 时,它会使用当前进程可见的所有 GPU:
from cudf_polars.engine.ray import RayEngine
with RayEngine() as engine:
result = query.collect(engine=engine)
这就是单节点多 GPU 执行所需的完整配置。不需要配置集群、不需要指定分区方案,也不需要修改查询本身。分块大小、溢写行为及回退模式等选项,通过所有流式引擎都接受的 StreamingOptions 对象设置。完整列表见配置选项参考。
引擎上层仍然是原样的 Polars。你编写同一个 LazyFrame,查询计划到达 cudf-polars 之前,Polars 优化器已经完成处理。我们在优化后的 IR 构建完成后接入,因此语义、类型推断和优化与 CPU 执行一致。GPU 引擎不支持的操作默认回退到 CPU 引擎。
RapidsMPF 如何实现这些能力
RapidsMPF 为 GPU 引擎提供两部分:一个流式执行框架,以及一组针对 GPU 间数据移动优化的通信原语。两者共同让查询突破单 GPU 的内存和计算能力限制。
流式执行与外存处理
引擎将输入分为数据块,使其流经查询图,逐块完成过滤、转换、聚合和连接。执行器会并发重叠运行多个操作,因此选择分块大小时,必须为各操作所需的中间缓冲区留出空间。
执行框架以 Hoare 的通信顺序进程(Communicating Sequential Processes)为模型。查询计划中的每个物理操作都成为一个长期运行的 actor 协程,actor 通过容量有界的通道连接:扫描为选择提供数据,选择再连接过滤,过滤最终连接输出端。有界通道提供背压,因此较慢的消费者会限制生产者,而不是在 GPU 内存中积累无界队列。即使同时有大量操作正在执行,内存使用仍然可以预测。
前面两项能力由此而来。内存紧张时,未在处理中的分块会溢写到主机内存,需要时再调回,因此单 GPU 能够处理规模为显存数倍的数据集。同样的分解方式也适用于多 GPU:让单 GPU 流式处理 1 TB 数据集的分块机制,也让八个 GPU 能够分担该数据集。
无论采用哪种引擎运行方式,都使用这一执行模型。RayEngine、DaskEngine 和 SPMDEngine 都驱动同一个流式执行器,差别仅在于 GPU worker 如何配置。因此,选择依据是偏好的部署方式,而不是性能。
优化的通信原语
跨 GPU 扩展数据帧引擎,在很大程度上是数据移动问题。连接、排序和高基数分组聚合都要求具有相同键的行最终位于同一个 GPU 上,因此需要全对全 shuffle。shuffle 是分布式查询引擎常见的故障点,因为每个参与进程——每个 GPU 对应一个进程,称为 rank——必须同时持有它发送的数据和接收的数据。
RapidsMPF 将 shuffle 实现为流式集合通信。每个 rank 在分区表的数据块产生时插入这些块;块一到达,RapidsMPF 就将其路由到拥有对应哈希键的 rank,因此为 shuffle 提供输入的操作还在生产更多数据时,数据已开始跨网络传输。所有 rank 完成插入后,每个 rank 提取现在归其所有的分块。AllGather 等其他集合通信也遵循同样模式。

各 rank 在分块产生时将其路由到所属 rank,并在继续发送的同时接收归其所有的分块。
只有看到全部输入后,才能知道某个 rank 的输出大小,因此 RapidsMPF 自行分配输出缓冲区。这也使它能够溢写这些缓冲区:内存紧张时,分块移入主机内存,等需要它们的操作恢复执行时再返回。由此,即使 shuffle 涉及的数据超过设备内存,也能完成。
集合通信下面是一层基于 UCX/UCXX 或 MPI 的 communicator 抽象,它统一处理 CPU 与 GPU 数据,让传输层选择任意两个 rank 之间适合的路径。因此,从单 GPU 扩展到多 GPU,改变的是 rank 数量,而非程序结构。
完整机制见 RapidsMPF 背景文档,其中详细介绍了 shuffle 架构、通道和 actor 模型。
基准测试
为了展示大规模运行效果,我们在 TB 级规模运行了 PDS-H 和 PDS-DS。注 1
在 SF1K(约 1 TB 未压缩数据)下,单 GPU 运行完整 PDS-H 测试套件的速度为 CPU 引擎的 3.2 倍,PDS-DS 为 2.2 倍:


规模因子 1000 的 PDS-H(左)与 PDS-DS(右):单个 B200 对比双路 Xeon。
在 SF3K(约 3 TB)下,扩展到八个 GPU 后,PDS-H 加速达到 23.2 倍,PDS-DS 达到 11.0 倍:


规模因子 3000 的 PDS-H(左)与 PDS-DS(右):八个 B200 对比双路 Xeon。
这些基准测试可以端到端复现。基准测试文档展示了使用 tpchgen-cli 生成 PDS-H 数据,以及 CPU、单 GPU 和多 GPU 运行器的调用方式。两个套件的查询实现均位于 cuDF 仓库。
用自己的数据试一试
如果你的 Polars 流水线运行时间过长,影响迭代,可以直接评估 GPU 引擎。安装 cudf-polars,并在现有 .collect() 调用中添加 engine="gpu",先获得单 GPU 基准。从这里开始,RayEngine 可将同一个查询分配到所有可用 GPU。不支持的操作回退到 CPU 引擎,因此可以原样运行现有查询。
Polars GPU 支持指南记录了当前限制;cudf-polars 文档深入介绍引擎、配置和内存调优。
引擎正在持续开发,欢迎反馈。功能请求和 API 覆盖缺口最好提交到 cuDF 仓库;也可以在 Polars Discord 找到我们。
注释
-
Polars Decision Support(PDS-H 和 PDS-DS)是由 TPC-H 和 TPC-DS 基准测试派生的开放实现。它们采用 TPC 数据模型、数据生成方法和查询工作负载,在不同数据集规模下衡量分析查询性能。虽然紧密遵循 TPC 规范,但它们并非经过官方审计或认证的 TPC 基准测试。因此,PDS 结果用于 PDS 内部的比较评估,不能与已发表的 TPC 基准测试结果直接比较。返回正文
原文:Polars GPU 引擎的新流式后端;作者:Brian Tepera(NVIDIA);发表于 2026-08-28。
版权归原作者及来源机构所有。版本、性能数字与活动时间保留原文语境;示例未在本环境执行。











暂无评论内容