You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python Polars基于Hadoop的低内存Parquet读写与处理实践问题

使用Polars低内存处理HDFS超大型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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 23:51:25