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

如何在Polars扫描S3 Parquet并关联时限制内存占用?

解决S3大型Parquet数据集与本地小DataFrame关联时的内存占用问题

问题场景

需要处理S3上按dt分区存储的大型Parquet数据集,核心逻辑为:

  1. 过滤指定日期范围内的数据
  2. 和本地小型sellers_df按seller_id做内关联
  3. 原代码使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 07:15:12