你让云端处理一大批数据,结果却大到 Python 一次装不下。过去,Polars Cloud 要么把结果写进存储系统,要么先完整写入临时存储,再让程序读回来。新发布的 0.10.0 增加了第三种办法:结果算出一批,就交回 Python 一批。程序可以更早动手,也不必把全部结果同时塞进内存。
这次更新不只打通了“最后一公里”。它还允许多条查询一起提交,引入实验性分布式规划器 Miso,并为本地部署实验性接入 HDFS。我们此前报道过,同一套 Polars 查询可从笔记本上的 9700 万行扩展到云端 160 亿行;0.10.0 进一步把重点转向云端执行层:结果怎样回来,多项任务怎样协同,数据怎样少搬一次。
本文事实均来自 Polars 官方博客,尚无第三方实测或交叉信源。
结果不用等到最后
新接口 sink_batches() 实现了结果流——像搬货时装满一车就先发走,而不是等整座仓库清空。分布式查询会把计算拆给多台机器;每当一批结果准备好,系统便用一个 DataFrame 调用用户提供的 Python 回调函数。回调可以做自定义处理,解决“结果太大,无法整体放进内存”的场景。
用户可用 chunk_size 控制每次回调前缓冲多少行。回调返回 True,还能提前终止查询。不过,灵活性带来了几条重要约束。
首先,Polars Cloud 中 maintain_order 默认是 False,不同 worker——也就是参与计算的工作进程——可以并行调用回调。若设为 True,回调才会串行执行。这个默认值与开源 Polars 的同名接口不同,不能想当然地套用旧经验。
其次,回调必须具备幂等性——同一批数据处理两次,结果仍应与处理一次相同。任务重试时,同一批次可能由不同 worker 重复送达。比如回调负责写入外部系统,就需要避免重复记账或重复插入。
最后,官方明确说,这条路径比 sink_parquet、sink_csv、sink_ipc 和 sink_iceberg 等原生写出方式慢得多。它补充的是自定义能力,不是替代原生 sink;博客没有给出具体性能数字。
多条查询开始一起算
0.10.0 还让分布式 pl.collect_all() 可以把多个 LazyFrame 作为一个查询提交。LazyFrame 可以理解为“先记下要做什么,暂不立即计算”的查询计划。每个 LazyFrame 都必须以自己的 sink 结尾,并设置 lazy=True,让写出动作进入整体计划;缺少 sink 的任务会被拒绝,因为系统没有一个明确结果可交付。
一起提交的价值在于共享工作。官方示例中,两条查询都读取同一份事件数据:一条筛选点击,另一条按用户汇总。系统会对合并后的计划执行“公共子计划消除”,也就是识别重复步骤,让共享扫描只读一次,而不是每条查询各读一遍。
Miso先来试运行
查询规划器负责把筛选、连接和汇总排成执行步骤,并决定数据如何在机器间移动。分布式环境里,数据搬运往往是计划的重要部分。
0.10.0 在现有的 naive 规划器之外加入实验性 Miso,可按查询用 planner="miso" 开启。参数还接受 auto 和 naive,但目前 auto 仍然选择 naive。官方预告 Miso 近期会成为默认规划器,不过这只是路线计划,不是已经发生的变化。现阶段,用户可以在 Dashboard 查看两种规划器生成的 stage graph——查询被拆成哪些执行阶段——并比较实际运行时间。
另一个减少搬运的变化针对 Hive 分区数据。连接和分组等操作会按键分配数据;如果扫描时已经识别出相同的分区方式,0.10.0 可以跳过一次 shuffle,也就是省去一轮机器间重新分发。官方称这会加速过去需要支付这次搬运成本的查询,但没有披露适用条件和效果数据。
云端路线伸向自建集群
对于 On-Prem——企业自行部署和管理的集群——0.10.0 实验性支持扫描 HDFS,并读取存放在 HDFS 上的 Iceberg 元数据。HDFS 是把文件分散保存在多台服务器上的文件系统;Iceberg 则用元数据管理大型分析表的结构、快照和文件清单,本身不是存储系统。
worker 通过纯 Rust 客户端访问 HDFS,不要求集群安装 JVM;数据访问发生在 worker,而不是发起查询的客户端。不过,这些能力默认关闭,需要额外的集群配置、依赖和存储参数,不能泛化为所有 Polars Cloud 场景都已稳定支持。
为什么值得关注
这几个功能连起来,显出一条清楚的路线:Polars 不再只关心单条查询能否从本地扩展到云端,还开始处理一个执行平台必须面对的问题——大结果怎样持续交付,多条查询怎样共享计算,规划器怎样减少跨机器搬运,以及自建数据基础设施怎样接入。
0.10.0 更像一次能力拼图,而非性能结论。它让 Python 能介入分布式结果处理,也把并发查询和规划优化摆上台面。但官方目前提供的是功能说明,不是基准测试。
局限与未知
- 全部信息来自 Polars 官方博客,缺少发布文档之外的交叉验证和第三方生产实测。
sink_batches()只被定性为明显慢于原生 sink;Miso 与跳过 shuffle 的实际收益都没有具体数据。- Miso、On-Prem HDFS 及 Iceberg 元数据路径仍属实验性功能,默认配置和稳定边界尚未完全展开。