PySpark与Polars基准测试:公平触发懒求值的合理方案探讨
公平对比PySpark与Polars的查询触发方案
核心原则:公平性+贴近业务
你的核心诉求是在超内存场景下,让两个框架都以懒求值的最优路径执行查询,避免因触发方式导致的性能偏差。当前用PySpark .collect() 完全不合理——它会把全量数据拉到Driver内存,既不符合超内存场景的业务逻辑,也和Polars的流式collect完全不对等,必须替换。
方案1:用write()触发(最贴近真实业务)
真实业务中,超内存数据的处理结果几乎都会写入存储(Parquet、CSV、数据库等),因此用write()触发执行是最贴合实际场景的选择:
- PySpark:使用分布式写入API,不会将数据拉到Driver内存,完全适配超内存场景:
# 写入Parquet(和业务常用格式一致) df.write.format("parquet").mode("overwrite").save("/path/to/output") - Polars:用对应格式的写入方法,默认会处理超内存数据的流式输出:
# 写入Parquet,自动适配超内存场景 df.write_parquet("/path/to/output.parquet", use_pyarrow=True)
这种方式下,两个框架都会按照业务实际的执行路径运行,测试结果直接反映真实业务的性能表现。
方案2:无输出触发(排除IO干扰,专注计算性能)
如果担心存储IO会影响性能测试的准确性,可以用无输出的action算子触发完整查询计划,同时避免内存过载:
- PySpark:用
foreachPartition触发分布式执行,不输出任何数据:# 触发全量查询执行,但不输出结果 df.rdd.foreachPartition(lambda partition: None) - Polars:用
sink或map_batches配合流式执行,不产生输出:# 流式执行全量查询,无输出 df.sink(lambda batch: None, streaming=True)
这种方式可以排除IO开销的干扰,更精准地对比两个框架的计算能力(CPU、内存占用)。
关键公平性细节
- 资源配置对齐:确保PySpark的Executor内存、核数,和Polars的线程数(
pl.set_threads(n))匹配,避免因资源差异导致的性能偏差。 - 测试数据集规模:必须使用超过单节点内存的数据集,才能触发两个框架的流式/分布式处理逻辑,否则测试结果不具备参考价值。
内容的提问来源于stack exchange,提问作者Ernest P W
相关产品推荐
相关产品推荐

