如何在Polars扫描S3 Parquet并关联时限制内存占用?
解决S3大型Parquet数据集与本地小DataFrame关联时的内存占用问题
问题场景
需要处理S3上按dt分区存储的大型Parquet数据集,核心逻辑为:
- 过滤指定日期范围内的数据
- 和本地小型
sellers_df按seller_id做内关联 - 原代码使用
scan_pyarrow_dataset加载数据集,但执行时Polars会先把整个日期范围的全量数据加载到内存,再执行关联,引发内存溢出问题。
原代码如下:
import pyarrow.dataset as ds import polars as pl import s3fs # S3 credentials secret = ... key = ... endpoint_url = ... # 本地小型DataFrame sellers_df = pl.DataFrame({'seller_id': ['0332649223', '0192491683', '0336435426']}) # 扫描、过滤并关联S3上的大型数据集 fs = s3fs.S3FileSystem(endpoint_url=endpoint_url, key=key, secret=secret) dataset = ds.dataset(f'{s3_bucket}/benchmark_dt/dt_partitions', filesystem=fs, partitioning='hive') scan_df = pl.scan_pyarrow_dataset(dataset) \ .filter(pl.col('dt') >= '2023-05-17') \ .filter(pl.col('dt') <= '2023-10-18') \ .join(sellers_df.lazy(), on='seller_id', how='inner').collect()
Parquet文件分区结构:
-- dt_partitions -- dt=2023-06-09 -- data.parquet -- dt=2023-06-10 -- data.parquet -- dt=2023-06-11 -- data.parquet -- dt=2023-06-12 -- data.parquet ...
问题诊断
原执行计划显示未启用流式处理,关联操作需等待全量数据加载完成后执行:
INNER JOIN: LEFT PLAN ON: [col("seller_id")] PYTHON SCAN PROJECT */3 COLUMNS SELECTION: ((pa.compute.field('dt') >= '2023-10-17') & (pa.compute.field('dt') <= '2023-10-18')) RIGHT PLAN ON: [col("seller_id")] DF ["seller_id"]; PROJECT */1 COLUMNS; SELECTION: "None" END INNER JOIN
若改用is_in替代join,过滤条件可下推到扫描阶段,但仍无法实现流式关联:
PYTHON SCAN PROJECT */3 COLUMNS SELECTION: ((pa.compute.field('seller_id')).isin(["0332649223","0192491683",...]) & ((pa.compute.field('dt') >= '2023-10-17') & (pa.compute.field('dt') <= '2023-10-18')))
有效解决方案
方案1:启用Polars原生S3支持,实现流式关联
添加环境变量允许HTTP协议(适配S3兼容存储),直接使用Polars的scan_parquet替代scan_pyarrow_dataset,Polars会自动实现流式处理,边扫描边关联,避免加载全量数据。
修改后的代码:
import os import polars as pl # 允许HTTP协议(适配S3兼容存储) os.environ['AWS_ALLOW_HTTP'] = 'true' # S3配置 os.environ['AWS_ACCESS_KEY_ID'] = key os.environ['AWS_SECRET_ACCESS_KEY'] = secret os.environ['AWS_ENDPOINT_URL'] = endpoint_url # 本地小型DataFrame sellers_df = pl.DataFrame({'seller_id': ['0332649223', '0192491683', '0336435426']}) # 流式扫描、过滤并关联 scan_df = pl.scan_parquet(f's3://{s3_bucket}/benchmark_dt/dt_partitions/**/*.parquet', hive_partitioning=True) \ .filter(pl.col('dt') >= '2023-05-17') \ .filter(pl.col('dt') <= '2023-10-18') \ .join(sellers_df.lazy(), on='seller_id', how='inner').collect()
修改后的执行计划显示流式处理已生效:
--- STREAMING INNER JOIN: LEFT PLAN ON: [col("seller_id")] Parquet SCAN s3://test-bucket/benchmark_dt/dt_partitions/dt=2023-10-17/part-0.parquet PROJECT */3 COLUMNS RIGHT PLAN ON: [col("seller_id")] DF ["seller_id"]; PROJECT */1 COLUMNS; SELECTION: "None" END INNER JOIN --- END STREAMING
方案2:用is_in提前过滤(简化场景)
如果无需保留关联逻辑的灵活性,可直接在过滤阶段用is_in筛选目标seller_id,将过滤条件下推到Parquet扫描,只加载符合条件的数据:
# 替换join为is_in过滤 scan_df = pl.scan_pyarrow_dataset(dataset) \ .filter(pl.col('dt').is_between('2023-05-17', '2023-10-18')) \ .filter(pl.col('seller_id').is_in(sellers_df['seller_id'])) \ .collect()
内容的提问来源于stack exchange,提问作者barak1412
相关产品推荐
相关产品推荐

