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

如何以内存高效方式从Polars LazyFrame采样?大Parquet数据场景

问题

处理CC NEWS Dataset的海量Parquet数据,受内存限制只能用pl.scan_parquet(path)加载为Polars LazyFrame,无法直接用pl.DataFrame加载。

尝试过以下代码:

import polars as pl

lf = pl.scan_parquet(path)
lf.select(pl.col("col_of_interest").sample(n=sample_size,seed=0)) \
   .sink_parquet("sample.parquet") # (1) 程序崩溃

# (2) 去掉.sink_parquet()后运行:
lf.collect() # 同样崩溃

目标列是爬取的网页内容大字段,只要涉及该列,哪怕sample_size=1程序也崩溃;但不包含该列时,LazyFrame能轻松转为DataFrame,小数据量列操作也正常。

求内存高效的LazyFrame采样方案,支持扫描/读取时采样的方法也可。

解决方案

方法1:读取阶段直接采样(内存友好度最高)

Parquet文件支持按行组读取,先获取文件的行组信息,随机选择部分行组加载后再采样,从根源避免加载全量数据:

import polars as pl
import random

# 获取Parquet文件的行组元数据
parquet_file = pl.scan_parquet(path)
row_groups = parquet_file.parquet_file_metadata.row_groups
total_row_groups = len(row_groups)

# 随机选择部分行组(示例选10%,可根据数据规模调整比例)
selected_row_groups = random.sample(range(total_row_groups), k=int(total_row_groups * 0.1))

# 仅加载选中的行组,再采样目标列
sampled_lf = (
    pl.scan_parquet(path, row_groups=selected_row_groups)
    .select("col_of_interest")
    .sample(n=sample_size, seed=0)
)

sampled_lf.sink_parquet("sample.parquet")

方法2:分块流式处理+累积采样

利用Polars的流式迭代功能,分块读取数据,每块采样少量数据,最后合并结果并二次采样到目标数量:

import polars as pl

target_sample_size = 100
per_chunk_sample = 10  # 每个数据块采样10条,可根据块大小调整

sampled_chunks = []
# 流式迭代读取数据块
for chunk in pl.scan_parquet(path).select("col_of_interest").streaming().iter_batches():
    chunk_sample = chunk.sample(n=per_chunk_sample, seed=0)
    sampled_chunks.append(chunk_sample)
    # 提前终止:累积样本足够时停止读取
    if sum(len(df) for df in sampled_chunks) >= target_sample_size:
        break

# 合并所有采样块,最终截取目标数量的样本
final_sample = pl.concat(sampled_chunks).sample(n=target_sample_size, seed=0)
final_sample.write_parquet("sample.parquet")

方法3:随机过滤实现近似采样

通过生成随机数过滤行,先筛选出小比例数据,再二次采样到目标数量,避免全量加载:

import polars as pl

# 设定初始采样比例(示例为0.001%,需根据数据总量调整)
sample_ratio = 0.00001

sampled_lf = (
    pl.scan_parquet(path)
    .select("col_of_interest")
    .filter(pl.random(seed=0) < sample_ratio)
)

# 如果初始采样结果过多,再二次采样到目标数量
sampled_lf = sampled_lf.sample(n=sample_size, seed=0)
sampled_lf.sink_parquet("sample.parquet")

关键注意点

原代码崩溃的核心原因:sample在LazyFrame中默认需要先将全量数据加载到内存才能执行,而大字段列内存占用极高,直接触发内存溢出。优先选择方法1,因为它从读取阶段就减少数据量,内存占用最低;若行组本身仍过大,再用方法2分块处理。

内容的提问来源于stack exchange,提问作者q.uijote

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 16:57:15