Polars 如何用 collect 的 engine 参数在内存、流式与 GPU 引擎间选择执行方式?

发布时间:2026/9/12 4:10:26
Polars 如何用 collect 的 engine 参数在内存、流式与 GPU 引擎间选择执行方式? Polars 如何用 collect 的 engine 参数在内存、流式与 GPU 引擎间选择执行方式【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars使用 Polars 的 Lazy API 时查询在调用collect前只是逐步累积到内部查询图中并不真正执行。真正决定数据用哪套引擎计算是在collect这一步默认走内存引擎传enginestreaming走流式引擎传enginegpu走 GPU 引擎需要 NVIDIA 硬件与 RAPIDS cuDF 后端Open Beta 阶段。本文基于仓库内的 执行文档、流式概念文档、GPU 支持文档 和 2.0 升级指南说明三种引擎各自的适用条件、切换方法和验证方式帮助你在写查询的最后一步做出选择。先理解 collect 的默认行为以文档中的 Reddit 数据集为例构建查询只是登记计算计划import polars as pl q ( pl.scan_csv(docs/assets/data/reddit.csv) .with_columns(pl.col(name).str.to_uppercase()) .filter(pl.col(comment_karma) 0) )调用collect后 Polars 才执行优化后的查询图。默认的内存引擎把全部数据当一个批次处理这意味着查询内存峰值时刻所有数据都要装进可用内存。文档给出的执行结果示例为 1000 万行中筛出shape: (14_029, 6)文档示例输出实际数值以你的数据为准df q.collect()判断依据很简单如果你的数据在查询内存峰值处放得下默认collect()就是最短路径。另外注意文档中的警告LazyFrame是查询计划每次在下游被复用都会重新计算且group_by这类不保持行序的操作每次运行行序可能变化需要时用maintain_orderTrue。内存放不下时切到流式引擎当数据需要的内存超过可用内存时把enginestreaming传给collectPolars 会分批batches执行查询从而处理放不进内存的数据集。流式文档同时说明除内存压力外流式引擎比内存引擎性能更好q ( pl.scan_csv(docs/assets/data/iris.csv) .filter(pl.col(sepal_length) 5) .group_by(species) .agg(pl.col(sepal_width).mean()) ) df q.collect(enginestreaming)两点边界需要注意部分算子会回退。一些操作天然不支持流式或尚未实现流式版本Polars 会对这些操作回退到内存引擎用户无需感知但排查内存或性能问题时有用。可用物理计划图检查。用show_graph以enginestreaming绘制流式查询的物理图图例标注了各操作可能有多大的内存开销q.show_graph(plan_stagephysical, enginestreaming)行序变化2.0 起engineautocollect的默认值对惰性查询解析为流式引擎而流式引擎对不要求行序的操作unpivot、group_by、join 等不保证行序。如果代码依赖了之前的行序要么显式排序要么在支持的操作上传maintain_order例如join(..., maintain_orderleft)。有 NVIDIA GPU 时启用 GPU 引擎GPU 引擎面向 Python 用户基于 RAPIDS cuDF处于 Open Beta 阶段。启用前有硬性前提来自 GPU 支持文档NVIDIA Volta 或更高 GPUcompute capability 7.0CUDA 12 或 CUDA 13Linux 或 Windows Subsystem for Linux 2WSL2安装 GPU 后端pip install polars[gpu]pip install polars[gpu]目前配置为安装 CUDA 12 对应的cudf-polars-cu12。如果你的系统 CUDA 版本不同需要单独安装带对应后缀的 cudf-polars 库例如 CUDA 13pip install polars cudf-polars-cu13文档提醒cudf-polars只支持一个有限的 Polars 版本区间。若不锁定版本包解析器可能选到较旧的兼容 Polars 发布版如果这种回退不可接受应把需要的 Polars 版本钉住不兼容的组合会让依赖解析直接失败。用 enginegpu 执行查询构建好惰性查询后把enginegpu传给.collect.sink_*同样接受该参数df pl.LazyFrame({a: [1.242, 1.535]}) q df.select(pl.col(a).round(1)) result q.collect(enginegpu) print(result)enginegpu适合单 GPU。多 GPU 执行和查询运行时配置通过传GPUEngine对象实现。文档说明 cudf-polars 26.06 版本提供 3 个GPUEngine子类RayEngine基于 Ray 的多 GPU 执行DaskEngine基于 Dask 的多 GPU 执行SPMDEngine单程序多数据模型的多 GPU 执行这些引擎会拉起资源可作为上下文管理器使用以便释放from cudf_polars.engine.ray import RayEngine with RayEngine() as engine: result q.collect(engineengine) print(result)注意使用RayEngine或DaskEngine需要分别安装 Ray 或 Dask可通过 cudf-polars 的[ray]或[dask]pip extra 安装例如 CUDA 13 下pip install cudf-polars-cu13[ray] pip install cudf-polars-cu13[dask]验证查询是否真的走了 GPUGPU 模式下遇到不支持的操作不会让查询失败而是透明回退到标准 Polars 引擎在 CPU 上执行因此执行时间可能没有任何变化。文档给出两种确认手段。一是开启 verbose 模式不能上 GPU 的查询会发出PerformanceWarningwith pl.Config() as cfg: cfg.set_verbose(True) result q.collect(enginegpu)文档示例中的告警输出文档示例具体原因随你的查询而定PerformanceWarning: Query execution with GPU not possible: unsupported operations The errors were: - NotImplementedError: dtypeBinary conversion not supported二是禁止回退让不支持的查询直接抛异常q.collect(enginepl.GPUEngine(raise_on_failTrue))不支持时得到polars.exceptions.ComputeError文档示例。文档同时说明目前只报告 GPU 执行失败的近因计划扩展为报告查询中所有不支持的操作。GPU 引擎的支持范围支持文档列举的高层类别LazyFrame API、SQL APICSV、Parquet、ndjson 和内存 CPU DataFrame 的 I/O数值、逻辑、字符串、日期时间类型的操作字符串处理聚合含分组与滚动变体、join、过滤、缺失数据、连接concatenation不支持Eager DataFrame APIDate、Categorical、Enum、Time、Array、Binary、Object 数据类型部分带时区 Datetime 与 List 类型表达式时间序列重采样、Folds、用户自定义函数Excel 和数据库文件格式另外两点机制说明GPU 执行只在 Lazy API 中可用查询执行结束后结果以普通 CPU 内存中的 Polars DataFrame 返回CPU 与 GPU 引擎都基于 Apache Arrow 列式内存格式数据可以在两者间快速移动一个引擎写的文件另一个引擎可以读。文档给出的选型经验是当工作负载以分组聚合和 join 为主时最可能看到 GPU 加速I/O 受限的查询 GPU 与 CPU 性能通常相近按其测试80GiB 显存的 GPU 可容纳约 1 TiB 原始数据集视工作负载而定。按数据与硬件选择引擎的决策路径三条判断依据都来自上述文档可以按顺序套用数据在查询峰值内存处放得下且依赖行序或内存引擎特性直接用默认collect()或显式collect(enginein-memory)。数据放不进内存或想让大查询分批执行collect(enginestreaming)并用show_graph(plan_stagephysical, enginestreaming)检查哪些算子内存开销大、是否有回退。有满足要求的 NVIDIA GPU且查询以分组聚合和 join 为主、只涉及支持类型collect(enginegpu)配合 verbose 告警或raise_on_failTrue确认没有静默回退。对于 2.0 用户升级指南 还给出了进程级的引擎偏好设置可以整体回到 1.x 的默认行为pl.Config.set_engine_affinity(in-memory) # 进程级 lf.collect(enginein-memory) # 单查询也可以设置环境变量POLARS_ENGINE_AFFINITYin-memory。SQL 同样受影响pl.sql(..., eagerTrue)在 2.0 中会走LazyFrame.collect()即流式引擎需要保持内存引擎时可用pl.sql(..., eagerFalse).collect(enginein-memory)。验证时collect的返回就是一个可检查的 DataFrame对比三种引擎下同一查询的结果内容即可确认行为一致行序差异除外需显式排序后比较GPU 路径再叠加 verbose 告警或raise_on_fail来判断是否真的在 GPU 上执行。【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考