如何以内存高效方式从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
相关产品推荐
相关产品推荐

