news 2026/9/12 4:10:20

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
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这一步:默认走内存引擎,传engine="streaming"走流式引擎,传engine="gpu"走 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_order=True

内存放不下时切到流式引擎

当数据需要的内存超过可用内存时,把engine="streaming"传给collect,Polars 会分批(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(engine="streaming")

两点边界需要注意:

  • 部分算子会回退。一些操作天然不支持流式或尚未实现流式版本,Polars 会对这些操作回退到内存引擎,用户无需感知,但排查内存或性能问题时有用。
  • 可用物理计划图检查。用show_graphengine="streaming"绘制流式查询的物理图,图例标注了各操作可能有多大的内存开销:
q.show_graph(plan_stage="physical", engine="streaming")

行序变化:2.0 起engine="auto"collect的默认值)对惰性查询解析为流式引擎,而流式引擎对不要求行序的操作(unpivotgroup_by、join 等)不保证行序。如果代码依赖了之前的行序,要么显式排序,要么在支持的操作上传maintain_order,例如join(..., maintain_order="left")

有 NVIDIA GPU 时启用 GPU 引擎

GPU 引擎面向 Python 用户,基于 RAPIDS cuDF,处于 Open Beta 阶段。启用前有硬性前提,来自 GPU 支持文档:

  • NVIDIA Volta 或更高 GPU,compute capability 7.0+
  • CUDA 12 或 CUDA 13
  • Linux 或 Windows Subsystem for Linux 2(WSL2)

安装 GPU 后端

pip install polars[gpu]

pip install polars[gpu]目前配置为安装 CUDA 12 对应的cudf-polars-cu12。如果你的系统 CUDA 版本不同,需要单独安装带对应后缀的 cudf-polars 库,例如 CUDA 13:

pip install polars cudf-polars-cu13

文档提醒:cudf-polars只支持一个有限的 Polars 版本区间。若不锁定版本,包解析器可能选到较旧的兼容 Polars 发布版;如果这种回退不可接受,应把需要的 Polars 版本钉住,不兼容的组合会让依赖解析直接失败。

用 engine="gpu" 执行查询

构建好惰性查询后,把engine="gpu"传给.collect.sink_*同样接受该参数):

df = pl.LazyFrame({"a": [1.242, 1.535]}) q = df.select(pl.col("a").round(1)) result = q.collect(engine="gpu") print(result)

engine="gpu"适合单 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(engine=engine) print(result)

注意:使用RayEngineDaskEngine需要分别安装 Ray 或 Dask,可通过 cudf-polars 的[ray][dask]pip extra 安装,例如 CUDA 13 下:

pip install cudf-polars-cu13[ray] pip install cudf-polars-cu13[dask]

验证查询是否真的走了 GPU

GPU 模式下遇到不支持的操作不会让查询失败,而是透明回退到标准 Polars 引擎在 CPU 上执行,因此执行时间可能没有任何变化。文档给出两种确认手段。

一是开启 verbose 模式,不能上 GPU 的查询会发出PerformanceWarning

with pl.Config() as cfg: cfg.set_verbose(True) result = q.collect(engine="gpu")

文档示例中的告警输出(文档示例,具体原因随你的查询而定):

PerformanceWarning: Query execution with GPU not possible: unsupported operations The errors were: - NotImplementedError: dtype=Binary conversion not supported

二是禁止回退,让不支持的查询直接抛异常:

q.collect(engine=pl.GPUEngine(raise_on_fail=True))

不支持时得到polars.exceptions.ComputeError(文档示例)。文档同时说明,目前只报告 GPU 执行失败的近因,计划扩展为报告查询中所有不支持的操作。

GPU 引擎的支持范围

支持(文档列举的高层类别):

  • LazyFrame API、SQL API
  • CSV、Parquet、ndjson 和内存 CPU DataFrame 的 I/O
  • 数值、逻辑、字符串、日期时间类型的操作,字符串处理
  • 聚合(含分组与滚动变体)、join、过滤、缺失数据、连接(concatenation)

不支持:

  • Eager DataFrame API
  • Date、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 原始数据集(视工作负载而定)。

按数据与硬件选择引擎的决策路径

三条判断依据都来自上述文档,可以按顺序套用:

  1. 数据在查询峰值内存处放得下,且依赖行序或内存引擎特性:直接用默认collect(),或显式collect(engine="in-memory")
  2. 数据放不进内存,或想让大查询分批执行collect(engine="streaming"),并用show_graph(plan_stage="physical", engine="streaming")检查哪些算子内存开销大、是否有回退。
  3. 有满足要求的 NVIDIA GPU,且查询以分组聚合和 join 为主、只涉及支持类型collect(engine="gpu"),配合 verbose 告警或raise_on_fail=True确认没有静默回退。

对于 2.0 用户,升级指南 还给出了进程级的引擎偏好设置,可以整体回到 1.x 的默认行为:

pl.Config.set_engine_affinity("in-memory") # 进程级 lf.collect(engine="in-memory") # 单查询

也可以设置环境变量POLARS_ENGINE_AFFINITY=in-memory。SQL 同样受影响:pl.sql(..., eager=True)在 2.0 中会走LazyFrame.collect()(即流式引擎),需要保持内存引擎时可用pl.sql(..., eager=False).collect(engine="in-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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 4:09:34

Android车载CAN开发:从SocketCAN到UDS诊断的全链路实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/12 4:08:32

Kilo CLI 贡献指南:从环境搭建到提交高质量 PR 的完整实战手册

Kilo CLI 贡献指南:从环境搭建到提交高质量 PR 的完整实战手册 【免费下载链接】kilocode Kilo is the all-in-one agentic engineering platform. Build, ship, and iterate faster with the most popular open source coding agent. 项目地址: https://gitcode.…

作者头像 李华
网站建设 2026/9/12 4:07:05

机器学习大作业实战指南:从线性回归到神经网络与不平衡分类

简介:一套面向西电机器学习课程的大作业资料包,适合计算机、电子信息、数学等专业学生完成期末大作业或课程设计时参考。内容覆盖10个实验,包括逻辑与二分类(C2-1)、带噪声的线性回归(C3-1)、神…

作者头像 李华