Polars 2.0
Polars 2.0
Today we are shipping Polars 2.0. In the earlier announcement post we went through the rationale of the version bump. This post we will discuss what features 2.0 brings. Even though we didn’t intend to make it a big feature release, it still packs a lot to get enthousiastic about. Let’s go through the highlights of this release: our initial version of out-of-core (spill-to-disk) support is enabled, a lot of very core performance improvements, first class SQL support, which together with the performance improvements has Polars leading DataFusion and DuckDB in TPC-H and TPC-DS1 benchmarks, a new Map dtype, and stricter Polars on dtypes and explicitness, leading to faster feedback, and faster AI iteration.
今天我们正式发布 Polars 2.0。在之前的公告文章中,我们已经解释了版本升级的缘由。本文将探讨 2.0 版本带来的新特性。尽管我们最初并未打算将其作为一个大型功能版本,但它依然包含许多令人兴奋的更新。让我们来看看本次发布的核心亮点:启用了初步的核外(out-of-core,即溢出到磁盘)支持;进行了大量核心性能优化;提供了顶级 SQL 支持,结合性能提升,使 Polars 在 TPC-H 和 TPC-DS1 基准测试中领先于 DataFusion 和 DuckDB;新增了 Map 数据类型;以及对数据类型和显式性提出了更严格的要求,从而带来更快的反馈和更高效的 AI 迭代。
Performance and SQL as a first class citizen
性能与顶级 SQL 支持
Polars 2.0 will be the marking point where we will treat SQL as first class citizen. Polars SQL coverage has increased dramatically over last few months. We know we have been building a solid engine for the last couple of years. In Polars 2.0, we want to enable that to more workloads, including SQL. To make this performant, we shipped many improvements to our optimizer and engine. The highlights here join reordering, much better common-subplan-elimination and dynamic predicates/bloom filters.
Polars 2.0 将成为我们将 SQL 视为“一等公民”的里程碑。过去几个月里,Polars 的 SQL 覆盖率大幅提升。我们深知过去几年一直在构建一个稳健的引擎,而在 Polars 2.0 中,我们希望将其扩展到更多工作负载,包括 SQL。为了实现高性能,我们对优化器和引擎进行了多项改进,重点包括连接重排序(join reordering)、更优的公共子计划消除(common-subplan-elimination)以及动态谓词/布隆过滤器(dynamic predicates/bloom filters)。
To see how we perform on typical SQL benchmarks, we ran Polars SQL on data derived from TPC-H and TPC-DS1 and ran it against the latest DuckDB release (1.5.6), DuckDB 2.0 alpha (2.0.0.dev2610011535) and the latest DataFusion release (54.0.0) on a c7a.4xlarge (16 vCPUs, 32GB RAM) and a c7a.metal (192 vCPUs, 384GB RAM). Every query ran 5 times in a hot setting, with a separate process per query and a 60 second timeout. The file cache was cleared between each engine/benchmark (not between queries). For every query we take the best of the 5 runs, and we compare engines on both the sum and the geometric mean of those query times.
为了评估在典型 SQL 基准测试中的表现,我们在 c7a.4xlarge(16 vCPU,32GB 内存)和 c7a.metal(192 vCPU,384GB 内存)机器上,使用源自 TPC-H 和 TPC-DS1 的数据,对比了 Polars SQL 与最新版 DuckDB (1.5.6)、DuckDB 2.0 alpha (2.0.0.dev2610011535) 以及最新版 DataFusion (54.0.0)。每个查询在热启动环境下运行 5 次,每个查询使用独立进程,并设置 60 秒超时。在每个引擎/基准测试之间(而非查询之间)会清除文件缓存。对于每个查询,我们取 5 次运行中的最佳成绩,并基于查询时间的总和及几何平均值来比较各引擎。
Polars and both DuckDB versions completed all queries. DataFusion timed out on TPC-DS q72 (and once on q67) and ran out of memory on TPC-H q18 on c7a.4xlarge; those queries are excluded from the results above for all engines. We observe that default Polars is fastest on all but one benchmarks. Polars has a constant overhead when we scale to 192 threads, which hurts small data queries. In fact we see that Polars limited to 32 cores is competitive or winning in all benchmarks. We have diagnosed the cause on our end and will hopefully fix this problem in the next release.
Polars 和两个版本的 DuckDB 均完成了所有查询。DataFusion 在 c7a.4xlarge 上运行 TPC-DS q72(以及一次 q67)时超时,并在 TPC-H q18 上内存溢出;这些查询在所有引擎的结果中均被剔除。我们观察到,除一个基准测试外,默认配置的 Polars 在其余测试中均为最快。当扩展到 192 个线程时,Polars 存在一定的固定开销,这会影响小数据查询的性能。事实上,我们发现将 Polars 限制在 32 核时,它在所有基准测试中都具有竞争力或处于领先地位。我们已经诊断出原因,并希望在下一个版本中解决此问题。
Streaming engine and OOC as default
流式引擎与默认开启的核外处理
This is the one of the biggest impact changes of 2.0. Calling collect on a LazyFrame will now default to the streaming engine, leading to massive memory and performance improvements on most queries. The reason this required a major version bump is that the streaming engine doesn’t guarantee row-order by default for certain operations (join, group_by, unpivot, etc.). If you require observable row-order in those operations, you can opt in to that by setting maintain_order=True.
这是 2.0 版本中影响最大的变化之一。现在,对 LazyFrame 调用 collect 将默认使用流式引擎,从而在大多数查询中带来巨大的内存和性能提升。之所以需要进行大版本升级,是因为流式引擎在某些操作(如 join、group_by、unpivot 等)中默认不保证行顺序。如果您在这些操作中需要可观察的行顺序,可以通过设置 maintain_order=True 来启用。
Out-of-core (spill to disk) is now enabled by default. It starts spilling at ~80% of RAM (this may need tuning). Operations that support out-of-core at this moment (sort, window functions, many expressions) can now start spilling to disk to finish a query. The default disk budget is 64GB. In the coming time we will enable out-of-core for joins and group-by’s as well. These two changes will make Polars much more resillient in high-memory workloads for casual data practicioners.
核外处理(溢出到磁盘)现已默认开启。当内存占用达到约 80% 时,系统将开始溢出(此阈值可能需要调整)。目前支持核外处理的操作(如排序、窗口函数及许多表达式)现在可以开始通过溢出到磁盘来完成查询。默认磁盘配额为 64GB。未来,我们还将为 join 和 group-by 操作启用核外处理。这两项改进将使 Polars 在处理高内存负载时,对于普通数据从业者而言更加稳健。
New Map datatype
新增 Map 数据类型
Polars now supports the Arrow MapType directly as a Polars Map dtype. You can think of a Map as a Python dictionary, mapping keys to values. Before 2.0 the Arrow MapType was read in Polars as List(Struct({“key”: …, “value”: …})).
Polars 现在直接支持将 Arrow MapType 作为 Polars 的 Map 数据类型。您可以将 Map 视为 Python 字典,即键值对映射。在 2.0 之前,Arrow MapType 在 Polars 中被读取为 List(Struct({"key": ..., "value": ...}))。