Python Polars基于Hadoop的低内存Parquet读写与处理实践问题
问题背景
需要用Polars处理HDFS上的超大型Parquet文件以避免内存溢出,官方文档推荐使用scan_parquet、LazyFrame和sink_parquet,但缺少HDFS环境下的实操指南。以下是内存内处理小文件的可行示例:
1. 环境配置
# 导入依赖 import pandas as pd import polars as pl import numpy as np import pyarrow as pa import pyarrow.parquet as pq from pyarrow._hdfs import HadoopFileSystem # 初始化HDFS文件系统 hdfs_filesystem = HDFSConnection('default') hdfs_out_path_1 = "scanexample.parquet" hdfs_out_path_2 = "scanexample2.parquet"
2. 生成数据并写入HDFS
# 生成测试数据集 df = pd.DataFrame({ 'A': np.arange(10000), 'B': np.arange(10000), 'C': np.arange(10000), 'D': np.arange(10000), }) # 将数据写入HDFS pq_table = pa.Table.from_pandas(df) pq_writer = pq.ParquetWriter(hdfs_out_path_1, schema=pq_table.schema, filesystem=hdfs_filesystem) # 写入数据并关闭写入器 pq_writer.write_table(pq_table) pq_writer.close()
3. 内存中读取HDFS上的Parquet文件
# 读取文件到Polars DataFrame pq_df = pl.read_parquet(source=hdfs_out_path_1, use_pyarrow=True, pyarrow_options={"filesystem": hdfs_filesystem})
4. 数据转换并写入HDFS
# 过滤数据并写入新文件 pq_df.filter(pl.col('A')>9000)\ .write_parquet(file = hdfs_out_path_2, use_pyarrow=True, pyarrow_options={"filesystem": hdfs_filesystem})
低内存处理时的报错情况
尝试用scan_parquet实现低内存处理时,出现以下错误:
尝试1:直接扫描HDFS路径
# 扫描文件:尝试1 scan_df = pl.scan_parquet(source = hdfs_out_path_2)
ERROR: Cannot find file
尝试2:传入HDFS文件流
# 扫描文件:尝试2 scan_df = pl.scan_parquet(source = hdfs_filesystem.open_input_stream(hdfs_out_path_1))
ERROR: expected str, bytes or os.PathLike object, not pyarrow.lib.NativeFile
Polars官方文档显示scan_parquet不支持pyarrow_options参数,但提及支持存储选项,不清楚具体配置方式。
本地无HDFS环境时,低内存处理流程可正常运行:
# 写入本地Parquet文件 df.to_parquet(path="testlocal.parquet") # 延迟读取生成LazyFrame lazy_df = pl.scan_parquet(source="testlocal.parquet") # 数据转换并低内存写入结果 lazy_df.filter(pl.col('A')>9000).sink_parquet(path= "testlocal.out.parquet")
更新:PyArrow Dataset加载后的问题
尝试用PyArrow Dataset加载为LazyFrame后,无法直接使用sink_parquet,必须先collect()到内存,不符合低内存需求:
# 读取为LazyFrame import pyarrow.dataset as ds pq_lf = pl.scan_pyarrow_dataset( ds.dataset(hdfs_out_path_1, filesystem= hdfs_filesystem)) # 尝试写入结果 pq_lf.filter(pl.col('A')>9000).sink_parquet(path= "testlocal.out.parquet")
PanicException: sink_parquet not yet supported in standard engine. Use 'collect().write_parquet()'
解决方案
方法1:通过storage_options配置HDFS
Polars的scan_parquet和sink_parquet支持通过storage_options参数传递已初始化的HDFS文件系统对象,实现低内存处理:
# 延迟扫描HDFS上的Parquet文件 lazy_df = pl.scan_parquet( source=hdfs_out_path_1, storage_options={"backend": "pyarrow", "fs": hdfs_filesystem} ) # 执行过滤并低内存写入HDFS lazy_df.filter(pl.col('A')>9000).sink_parquet( file=hdfs_out_path_2, storage_options={"backend": "pyarrow", "fs": hdfs_filesystem} )
方法2:PyArrow Dataset逐批次处理
若方法1存在兼容问题,可通过PyArrow Dataset流式读取,配合Polars逐批次处理,最后用PyArrow写入HDFS:
import pyarrow.dataset as ds # 创建PyArrow Dataset dataset = ds.dataset(hdfs_out_path_1, filesystem=hdfs_filesystem) # 定义单批次处理逻辑 def process_batch(batch): pl_df = pl.from_arrow(batch) return pl_df.filter(pl.col('A')>9000).to_arrow() # 流式处理并写入结果 writer = None for batch in dataset.to_batches(): processed_batch = process_batch(batch) if writer is None: writer = pq.ParquetWriter(hdfs_out_path_2, schema=processed_batch.schema, filesystem=hdfs_filesystem) writer.write_table(processed_batch) if writer: writer.close()
内容的提问来源于stack exchange,提问作者Esben Eickhardt

