如何拆分超大列表以适配Polars的explode()内存限制?
解决Polars超大列表
explode的内存限制问题 问题背景
当使用Polars对包含超大列表(例如单列表含30万元素,对应字符串长度达670万)的列执行explode操作并流式写入文件时,会触发内存不足问题——Polars当前架构要求单行展开后的结果必须能完整放入内存。需要将超大列表拆分为固定大小的子列表(比如每5万元素一组),把1行数据拆分为多行小列表后再执行explode,以适配流式处理的内存限制。
原问题代码(会触发内存问题):
import polars as pl test = pl.LazyFrame({'col1': 'string '*1_000_000}) (test .with_columns(explode_me = pl.col('col1').str.split(' ')) .explode(pl.col('explode_me')) .collect(streaming=True) .write_parquet('file.parquet') )
可行解决方案
通过将大列表拆分为嵌套的小列表,先explode外层得到多行小列表,再直接流式写入文件,避免全量加载数据到内存。
优化实现代码:
import polars as pl test = pl.LazyFrame({'col1': 'string '*1_000_000}) (test .with_columns(explode_me = pl.col('col1').str.split(' ')) .with_columns( # 按每10000个元素拆分大列表为子列表,可根据内存调整步长 pl.col('explode_me').map_elements( lambda x: [x[i:i+10_000] for i in range(0, len(x), 10_000)], return_dtype=pl.List(pl.List(pl.String)) ) ) .explode('explode_me') # 展开外层嵌套,得到每行对应一个小列表 .sink_parquet('file.parquet') # 直接流式写入,无需全量collect )
关键细节说明
- 拆分逻辑:利用
map_elements遍历每个大列表,按指定步长切割成多个子列表,将原列的List(String)类型转换为List(List(String))的嵌套结构。 - 流式写入优化:使用
sink_parquet替代collect(streaming=True).write_parquet,在流式处理过程中直接写入文件,进一步降低内存占用。 - 步长灵活调整:可根据自身内存情况调整子列表的大小(比如改为50000),平衡内存占用和处理的行数。
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

