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

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进程相同的进程生成,但太快无法看清细节。

我想了解:

  1. 是否有比当前循环/映射更好的处理方法;
  2. 可进行哪些内存管理优化;
  3. 是否单纯需要更多资源。

更新内容

我修改了代码,不再创建新列后删除旧列,而是直接修改列(类似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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 20:15:32