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

Dask处理大Parquet文件遇段错误与CPU利用率低问题求助

问题背景

我有一个18GB的.parquet文件,包含约3亿条会计数据(无法共享),分为53个行组。任务是仅保留每个单元格字典内的特定词汇来清洗数据。文件读取正常,但在20核128GB内存的Ubuntu 22.04桌面处理时出现段错误。

使用Python和Dask库转换为包含以下列的Dask DataFrame:

['rowid', 'txid', 'debit', 'credit', 'effective_date', 'entered_date', 'user_id', 'transaction', 'memo', 'type', 'account', 'total_amt']

需要清洗的列为memo、type和account,当前处理方式是对这些列分别应用filter_field和hash_field方法:

if isinstance(data, dd.DataFrame):
    # clean memo columns
    # data = data.repartition(npartitions=20) <incl. this line in experiments with partition size>
    result = [data[col].apply(lambda x: self.filter_field(text=x, word_dict=word_dict), meta=data[col]) for col in memo_columns]
    for i, col in enumerate(memo_columns):  # this loop seems to be req'd to assign values
        data[col] = result[i]
    
    # second: hash the name/id fields
    id_cols = name_columns + id_columns + account_columns
    result = [data[col].apply(lambda x: self.hash_field(text=x), meta=data[col]) for col in id_cols]
    for i, col in enumerate(id_cols):
        data[col] = result[i]
        
    del result
    gc.collect()

filter_field会移除单元格中的符号,保留字典内的词汇;hash_field逻辑为shake_256(text.encode(encoding='utf8')).hexdigest(20)。

处理1.28亿行以内的同类文件正常,但处理大文件时出现两个问题:

  • 随机出现段错误,时间不定;
  • 仅3个左右核心利用率达30%-50%,其余核心仅约1%-2%(通过htop观察)。

想了解:

  1. 是否有比当前循环/矢量化更好的处理方法?
  2. 如何让更多核心参与处理?

补充说明:

  • 尝试过调整分区数量和大小,无明显改善,仍会出现段错误;
  • 预期处理逻辑是将各列拆分为约20个分片,并行处理(每个分片对应一个处理器)。

解决方案

1. 优化处理逻辑,替代逐行循环

(1)向量化重构filter_field

原apply逐行调用Python函数的方式效率极低,还容易引发内存问题。可以用正则向量化实现匹配:

import re
import dask.dataframe as dd

# 预编译正则表达式,提前转义特殊字符避免匹配异常
pattern = re.compile(r'\b(' + '|'.join(re.escape(word) for word in word_dict) + r')\b')

def vectorized_filter(series):
    # 批量提取匹配词汇并拼接
    return series.str.findall(pattern).str.join(' ')

# 批量处理目标列,直接赋值避免中间结果堆积
data[memo_columns] = data[memo_columns].apply(vectorized_filter, meta=str)

这种方式利用Dask的向量化字符串操作,减少GIL锁竞争,同时降低内存开销。

(2)向量化重构hash_field

同样避免逐行哈希,改用Dask字符串操作批量处理:

from hashlib import shake_256

def vectorized_hash(series):
    return series.str.encode('utf-8').apply(lambda x: shake_256(x).hexdigest(20), meta=str)

# 批量处理哈希列
data[id_cols] = data[id_cols].apply(vectorized_hash, meta=str)

(3)简化中间结果流程

原代码先创建result列表再循环赋值,会产生多余的Dask中间对象。直接对列批量赋值,减少内存中挂起的任务数量,降低段错误概率。

2. 提升核心利用率,解决并行不足

(1)切换进程调度器绕开GIL限制

默认线程调度器会受Python GIL锁限制,无法真正利用多核心。改用进程调度器:

from dask.distributed import Client

# 启动本地集群,预留2个核心给系统,每个worker分配6GB内存
client = Client(n_workers=18, threads_per_worker=1, memory_limit='6GB')

每个worker是独立进程,不受GIL约束,能真正并行处理任务。

(2)优化分区策略匹配核心数

分区数建议设置为核心数的2-4倍(比如20核对应40-80个分区),让Dask能更均匀地调度任务:

# 重新分区为核心数的3倍,同时控制单分区大小在100-200MB
data = data.repartition(npartitions=60)
# 或者读取时直接指定分片大小,避免后续重分区开销
data = dd.read_parquet('your_file.parquet', chunksize='200MB')

(3)排查段错误根源

段错误通常和内存溢出或底层库bug有关:

  • 给filter_field和hash_field增加空值处理逻辑,比如判断text是否为None再执行操作;
  • 通过Dask Dashboard监控worker内存,若出现内存超限,调整memory_limit并开启worker自动重启;
  • 升级Dask、PyArrow到最新版本,修复已知的底层bug。

内容的提问来源于stack exchange,提问作者MikeB2019x

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 03:49:56