Polars Parquet 读写与扫描实战指南:从read_parquet到惰性查询优化
【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars
本指南对应仓库
docs/source/user-guide/io/parquet.md,以 Python / Rust 双语言示例贯穿 Parquet 的读取、写入与惰性扫描三大主题。Polars 以 Rust 实现的高性能列式引擎,其内存中的DataFrame布局与磁盘上的 Parquet 列式文件布局高度相似,因此读写 Parquet 非常高效。读完本文,你将掌握read_parquet/scan_parquet/write_parquet的完整用法与关键参数,理解扫描为何比直接读取更适合大数据量与云上场景,并能在工程中正确选择"读"与"扫"。
Parquet 为什么与 Polars 如此契合
Parquet 与 CSV 这类面向行的文本格式存在本质差异:Parquet 是列式存储格式(columnar format)。数据按"列"而非按"行"组织在文件中,这种布局带来两方面的直接收益:
- 更好的压缩率:同一列中的数据往往具有相似的类型与取值分布,压缩算法可以对同构数据发挥更强作用;
- 更快的数据访问:当查询只关心少数几列时,引擎无需读取无关列的数据。
更进一步,Polars 之所以读写 Parquet 快,文档给出的核心解释是:PolarsDataFrame在内存中的布局,与 Parquet 文件在磁盘上的布局在许多方面是镜像对应的。列式内存布局 + 列式磁盘文件,意味着数据几乎无需做破坏性重排即可完成编解码与转换,这正是 Parquet 在 Polars 生态中成为默认高性能交换格式的根本原因。
在仓库中,这一能力的底层支撑分布在两个层面:Python 侧对外 API 位于 py-polars/src/polars/io/parquet/functions.py;Rust 侧则集中在 crates/polars-io/src/parquet/mod.rs 及其read/、write/子目录,并借助crates/polars-parquet/完成 Arrow 数组与 Parquet 页数据之间的转换。
读取:read_parquet
将本地 Parquet 文件读入DataFrame,直接调用read_parquet即可:
=== ":fontawesome-brands-python: Python"
import polars as pl df = pl.read_parquet("docs/assets/data/path.parquet")完整可运行示例见 docs/source/src/python/user-guide/io/parquet.py。
=== ":fontawesome-brands-rust: Rust"
use polars::prelude::*; fn main() -> Result<(), Box<dyn std::error::Error>> { let mut file = std::fs::File::open("docs/assets/data/path.parquet")?; let df = ParquetReader::new(&mut file).finish()?; println!("{df}"); Ok(()) }Rust 侧对应示例完整代码见 docs/source/src/rust/user-guide/io/parquet.rs。
常用参数速查
从 functions.py 的签名看,read_parquet在工程实战中常用的参数包括:
| 参数 | 默认值 | 说明 |
|---|---|---|
source | 必填 | 本地/云上路径,支持 glob;也接受具有read()方法的文件对象(如open()句柄、BytesIO) |
columns | None | 只读取指定列,接受列名列表或从 0 起始的列索引列表 |
n_rows | None | 仅读取前n_rows行;仅在use_pyarrow=False时有效 |
parallel | "auto" | 并行策略,可选'auto'、'columns'、'row_groups'、'none',auto自动选择最优方向 |
use_statistics | True | 利用文件页脚中的统计信息(如 min/max)跳过无需读取的数据页 |
hive_partitioning | None | 是否从 Hive 风格分区路径推断并裁剪;传入单个目录时自动启用,否则默认关闭 |
try_parse_hive_dates | True | 是否尝试把 Hive 分区值解析为 date/datetime |
low_memory | False | 以部分性能换更低内存占用 |
use_pyarrow | False | 切换为 PyArrow 读取器(官方描述为"更稳定"),代价是失去部分原生功能 |
storage_options | None | 云存储连接配置(AWS/GCP/Azure/Hugging Face),缺省时尝试从环境变量推断 |
credential_provider | "auto" | 提供云凭证的函数或CredentialProvider*工具类 |
memory_map | True | 内存映射底层文件,通常提升性能;仅在use_pyarrow=True时使用 |
row_index_name/row_index_offset | None/0 | 在结果首列插入行号列,可指定起始偏移 |
值得注意的兼容性细节(源码中已用RenamedParameter标注):0.20.4 起row_count_name/row_count_offset更名为row_index_name/row_index_offset,旧名在 2.0 中被移除;同样被移除的还有rechunk、retries(后者需改为通过storage_options={"max_retries": n}传递)。
进阶辅助函数
若只想获取元数据而不加载数据,同文件 还提供了两个轻量入口:
read_parquet_schema(source):仅返回列名到数据类型的 schema 字典(内部实现即scan_parquet(source).collect_schema());read_parquet_metadata(source):读取文件级自定义元数据字典(该 API 标注为实验性)。
写入:write_parquet
写入 Parquet 与读取同样直观:先在内存中构建一个DataFrame,再调用write_parquet指定目标路径:
=== ":fontawesome-brands-python: Python"
import polars as pl df = pl.DataFrame({"foo": [1, 2, 3], "bar": [None, "bak", "baz"]}) df.write_parquet("docs/assets/data/path.parquet")=== ":fontawesome-brands-rust: Rust"
use polars::prelude::*; fn main() -> Result<(), Box<dyn std::error::Error>> { let mut df = df!( "foo" => &[1, 2, 3], "bar" => &[None, Some("bak"), Some("baz")], )?; let mut file = std::fs::File::create("docs/assets/data/path.parquet")?; ParquetWriter::new(&mut file).finish(&mut df)?; Ok(()) }上例中的示例数据刻意包含了None空值列,用以说明 Parquet 对缺失值的表达能力——空值被编码在独立的有效性位图中,不会像 CSV 那样引入解析歧义。
从 Rust 代码可见,写入的底层入口是ParquetWriter::new(&mut file).finish(&mut df),对应实现位于 crates/polars-io/src/parquet/write/。此外,与本地路径一致,write_parquet同样支持传入s3://...等云存储 URL 直接上云(详见下文"云端场景")。
扫描:scan_parquet与惰性执行
与read_parquet立即解析不同,scan_parquet返回的是一个惰性计算持有者LazyFrame:调用扫描时文件的实际解析并不会立刻发生,查询计划被推迟到collect时才真正执行。
=== ":fontawesome-brands-python: Python"
import polars as pl df = pl.scan_parquet("docs/assets/data/path.parquet")=== ":fontawesome-brands-rust: Rust"
use polars::prelude::*; fn main() -> Result<(), Box<dyn std::error::Error>> { let args = ScanArgsParquet::default(); let lf = LazyFrame::scan_parquet(PlRefPath::new("docs/assets/data/path.parquet"), args)?; println!("{}", lf.collect()?); Ok(()) }为什么扫描更值得推荐
扫描的价值在于:LazyFrame允许查询优化器在真正读取数据之前做全局优化。Polars 的核心优化手段——谓词下推(predicate pushdown)与投影下推(projection pushdown)——能够被下推到扫描层,从而让文件读取阶段就只取出真正需要的行与列,典型效果是既提升速度又降低内存占用。
关于这些优化为何值得期待的完整解释,参见 用户指南 · 惰性 API 概念 与 优化项说明。
反模式警示:read_parquet().lazy()
一个在源码文档字符串中被明确点名的反模式是:
# 反模式:先物化整个文件再转惰性,无法把优化下推进读取器 df = pl.read_parquet("path.parquet").lazy()由于read_parquet已把文件完整物化为 eagerDataFrame,后续任何谓词/投影都无法再推入读取层。官方建议:凡是最终要以LazyFrame工作的场景,一律直接使用scan_parquet。
事实上,从实现上可以印证这一点:在use_pyarrow=False的默认路径下,read_parquet内部就是先构造scan_parquet(...)再立即执行_collect_eager()(见 functions.py)。换句话说,eager 读 = 扫描后立刻收集。
扫描的关键优化参数
scan_parquet签名 中,除与read_parquet重合的参数外,以下参数直接影响扫描性能:
| 参数 | 默认值 | 说明 |
|---|---|---|
use_statistics | True | 使用文件统计信息裁剪需读取的数据页,跳过不满足谓词的行组/页 |
parallel | "auto" | 除'auto'、'columns'、'row_groups'、'none'外,还可选'prefiltered'策略 |
glob | True | 是否对路径执行 glob 规则展开(一次扫描多个分片文件) |
cache | True | 是否缓存扫描结果,供同一逻辑计划中的多处使用复用 |
hive_partitioning | None | 是否启用 Hive 分区推断(传入单个目录时自动开启) |
hidden_file_prefix | None | 指定作为隐藏文件前缀的字符串,用于扫描目录时过滤文件 |
low_memory | False | 降低内存压力换取部分性能 |
其中parallel='prefiltered'是一个值得在大型文件上尝试的新策略:它先在并行条件下评估下推后的谓词,得到"哪些行需要读取"的掩码;随后再对列与行组同时并行、并在读取时过滤掉无需的行。官方注释给出的适用判断是:对于含大量行组 + 谓词能明显过滤聚集行或过滤比例高的文件,prefiltered可能带来显著加速;反之则可能拖慢扫描。并且当没有谓词可下推时,该策略会自动回退到auto。这些说明均可在 functions.py 的parallel参数文档中查到。
云端扫描:把优化推进到数据下载之前
当 Parquet 存放在云对象存储时,扫描的优势会被进一步放大。以 S3 为例:
import polars as pl source = "s3://bucket/*.parquet" df = ( pl.scan_parquet(source) .filter(pl.col("id") < 100) .select("id", "value") .collect() )由于谓词与投影被下推进scan_parquet,Polars 会先裁剪出真正需要的字节范围再发起下载,从而显著减少网络传输量;真正的查询求值由collect()触发。这就是 cloud-storage.md 中 "Scanning from cloud storage with query optimisation" 一节所强调的核心收益——若改用 eager 读取,整份文件必须先被下载到本地,云端的网络优势将荡然无存。
云端场景还支持:
- 认证配置:通过
storage_options传入访问密钥(如aws_access_key_id、aws_secret_access_key、aws_region);或用pl.CredentialProviderAWS等工具类选择 profile / 承担 IAM 角色;也可以自定义返回凭证字典与过期时间的函数,并通过pl.Config.set_default_credential_provider(...)设为全局默认; - 重试配置:
storage_options支持max_retries、retry_init_backoff_ms、retry_max_backoff_ms、retry_timeout_ms等键,精细化控制重试与退避行为; - PyArrow 数据集扫描:
scan_pyarrow_dataset(ds.dataset("s3://...", format="parquet"))适合 Hive 分区等复杂数据集(该功能依赖 PyArrow)。
完整的云端读写代码清单见 docs/source/src/python/user-guide/io/cloud-storage.py。
何时选择读、写还是扫
把三种操作放在一起对照,决策就非常清晰:
| 操作 | 返回类型 | 行为 | 适用场景 |
|---|---|---|---|
read_parquet | DataFrame | 立即解析并物化 | 小文件、单次性分析、交互式探索 |
write_parquet | 无(写文件) | 将内存中的DataFrame落盘/上云 | 结果持久化、ETL 落盘、构建分析数据集 |
scan_parquet | LazyFrame | 延迟解析,构建惰性查询计划 | 大文件、多文件 glob、云端读取、需要谓词/投影下推的复杂查询链 |
小结
Polars 之所以将 Parquet 视为一等公民格式,根源在于两种列式布局的高度同构;而scan_parquet又把这份效率优势延伸到惰性查询优化与云端场景中。实践中请记住三条原则:需要懒执行就用scan_parquet而非read_parquet().lazy();大文件与多分片场景优先扫描并利用use_statistics、hive_partitioning、glob 与prefiltered并行策略;云上数据务必让谓词与投影下推进读取层,让下载量只覆盖查询真正需要的字节。
【免费下载链接】polarsExtremely fast Query Engine for DataFrames, written in Rust项目地址: https://gitcode.com/GitHub_Trending/po/polars
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考