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观察)。
想了解:
- 是否有比当前循环/矢量化更好的处理方法?
- 如何让更多核心参与处理?
补充说明:
- 尝试过调整分区数量和大小,无明显改善,仍会出现段错误;
- 预期处理逻辑是将各列拆分为约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

