Polars处理7GB Parquet文件遇段错误及写入异常求助
我有一个7GB的.parquet文件,包含1.28亿行会计数据(无法共享),分为53个行组。我的任务是通过保留每个单元格中字典内的特定词汇来清洗数据。文件读取无问题,但在20核、128GB内存的Ubuntu桌面系统上处理时出现段错误。
我使用Python的Polars库将数据转换为包含以下列的DataFrame:
['rowid', 'txid', 'debit', 'credit', 'effective_date', 'entered_date', 'user_id', 'transaction', 'memo', 'type', 'account', 'total_amt']
需要清洗的列是memo、type和account。我的做法是遍历这些列并应用filter_field方法,问题出在处理循环中:
# two-step cleaning: first memo/account fields for memo_field in memo_columns+account_columns: print('memo_field:', memo_field) data = data.with_columns( (pl.col(memo_field).map_elements(lambda x: self.filter_field(text=x, word_dict=word_dict))).alias('clean_' + memo_field) # baseline ) data = data.drop(memo_field) #drop cleaned column
每次循环中,我会创建新的“clean”列并删除原列。
我确认过滤逻辑是可靠的,因为处理最多7800万行的同类文件时一切正常。处理更大的文件时,内存消耗持续攀升直至出现段错误,且段错误前在htop中看到十几个与主Python进程相同的进程生成,但太快无法看清细节。
我想了解:
- 是否有比当前循环/映射更好的处理方法;
- 可进行哪些内存管理优化;
- 是否单纯需要更多资源。
更新内容
我修改了代码,不再创建新列后删除旧列,而是直接修改列(类似Pandas的inplace操作)。这使得整个文件可以被处理,但每列清洗耗时约1500秒。
然而,现在出现了写入错误:Parquet cannot store strings with size 2GB or more。这让我困惑,因为清洗后数据量应该只会减少。
1. 替代循环/映射的更优处理方式
- 弃用
map_elements,改用Polars原生矢量化表达式:map_elements是逐行调用Python函数,会触发GIL锁,无法利用多核优势,性能极低。如果你的filter_field逻辑是保留字典内的词汇,用Polars原生字符串方法重写即可实现矢量化执行:
假设word_dict是集合类型,要保留单元格中属于该集合的词汇并拼接,代码示例:
原生表达式性能比word_set = pl.lit(list(word_dict)) # 批量处理所有目标列,无需循环 data = data.with_columns( pl.col(memo_columns + account_columns) .str.split(" ") # 根据实际文本分隔符调整 .list.set_intersection(word_set) .list.join(" ") .alias(lambda col: f"clean_{col}") ).drop(memo_columns + account_columns)map_elements高几个数量级,还能避免Python层的额外内存开销。 - 批量处理列:一次性处理所有需要清洗的列,减少中间DataFrame副本的生成次数。
2. 内存管理优化
- 使用Polars Lazy API流式处理:不要一次性加载全量数据到内存,改用
scan_parquet生成LazyFrame,Polars会自动分块处理并优化执行计划,大幅降低内存占用:# 以流式方式读取文件 data = pl.scan_parquet("your_file.parquet") # 应用清洗逻辑 data = data.with_columns( pl.col(memo_columns + account_columns) .str.split(" ") .list.set_intersection(pl.lit(list(word_dict))) .list.join(" ") .alias(lambda col: f"clean_{col}") ).drop(memo_columns + account_columns) # 直接写入结果,无需加载全量数据到内存 data.sink_parquet("cleaned_file.parquet") - 避免不必要的中间副本:原循环中每次
with_columns+drop会生成多个DataFrame副本,内存持续累积。用Lazy API或一次性处理所有列可彻底避免这个问题。 - 开启流式处理模式:强制Polars以流式方式处理数据,进一步限制内存使用:
pl.activate_streaming() - 检查
filter_field内存泄漏:如果自定义函数中存在全局变量或未及时释放的对象,可能导致内存泄漏。改用原生表达式后可彻底规避这类问题。
3. 是否需要更多资源?
20核128GB内存的配置完全足够处理7GB的Parquet文件,问题出在处理方式而非资源不足。通过上述优化后,内存占用会大幅降低,性能也会显著提升,无需额外增加资源。
解决写入错误:Parquet cannot store strings with size 2GB or more
这个错误说明清洗后某一行的字符串长度超过了Parquet的单值大小限制(2GB),大概率是清洗逻辑在某些异常行生成了超长字符串(比如分隔符错误导致拼接大量内容,或空值/特殊字符处理不当)。
解决方法:
- 先排查异常行:用Lazy API筛选出目标列中字符串长度异常的行:
data = pl.scan_parquet("your_file.parquet") length_check = data.select( pl.col(memo_columns + account_columns) .str.lengths() .max() ).collect() print(length_check) - 添加长度限制:在清洗逻辑中限制字符串最大长度,避免触发Parquet的限制:
data = data.with_columns( pl.col(memo_columns + account_columns) .str.split(" ") .list.set_intersection(pl.lit(list(word_dict))) .list.join(" ") .str.slice(0, 1048576) # 限制最大长度为1MB .alias(lambda col: f"clean_{col}") ) - 预处理异常行:提前过滤或处理超长、空白、包含特殊字符的行:
data = data.filter( pl.col(memo_columns + account_columns).str.lengths() < 1000000 # 过滤超长行 )
内容的提问来源于stack exchange,提问作者MikeB2019x

