Python多线程.parquet单元格编码检测脚本CPU利用率低的优化问题
问题根源分析
- 任务粒度太小:每个单元格作为独立进程任务,进程间数据传递(IPC)的开销远大于编码检测的计算开销,导致CPU大部分时间消耗在数据传输上,无法有效利用。
- 单文件串行处理:脚本逐个读取Parquet文件、生成任务列表后再提交到进程池,进程池在等待文件读取、任务列表生成时处于闲置状态。
- 低效的文件IO:每个成功解码的单元格都创建独立JSON文件,还需循环判断文件名是否存在,大量小文件IO会严重阻塞进程,拖慢整体速度。
- Pandas逐行遍历低效:
df.iterrows()是逐行遍历DataFrame的低效方式,生成任务列表的过程耗时久,且这段时间CPU无法并行工作。
优化方案
1. 调整任务粒度:按文件/分块处理
将整个文件或文件分块作为单个任务提交到进程池,大幅减少IPC开销,让CPU专注于编码检测计算。
示例代码:
import os import pandas as pd import cchardet as chardet import json from multiprocessing import Pool def process_file(file_path, read_path): df = pd.read_parquet(file_path) results = [] # 按列遍历,比逐行更高效 for column in df.columns: series = df[column] for idx, value in series.items(): try: detect_result = chardet.detect(value) encoding = detect_result['encoding'] if encoding: text = value.decode(encoding) results.append({ "text": text, "encoding": encoding, "row_index": idx, "column_name": column }) # 减少控制台打印,可改为日志记录 # print(f"Text: {text}, Encoding: {encoding}") except (UnicodeDecodeError, LookupError): continue # 每个文件生成一个结果JSON,避免小文件泛滥 if results: output_filename = os.path.basename(file_path).replace('.parquet', '_results.json') with open(os.path.join(read_path, output_filename), 'w') as f: json.dump(results, f, indent=2) def main(): folder_path = '/home/master/Documents/smb/' read_path = '/home/master/Documents/smb/read/' # 一次性获取所有Parquet文件路径 file_paths = [ os.path.join(folder_path, filename) for filename in os.listdir(folder_path) if filename.endswith('.parquet') ] # 进程池一次性提交所有文件任务,充分利用CPU核心 with Pool(processes=os.cpu_count()) as pool: pool.starmap(process_file, [(path, read_path) for path in file_paths]) pool.close() pool.join() if __name__ == '__main__': main()
2. 合并文件IO操作
放弃每个单元格生成独立JSON的逻辑,改为按文件或分块生成单个结果文件,彻底消除小文件IO带来的阻塞,释放CPU资源。
3. 优化大文件处理:分块读取
针对超大Parquet文件,使用chunksize参数分块读取,避免一次性加载全量数据到内存,同时将分块作为任务提交,进一步提升并行效率。
示例分块处理逻辑:
def process_chunk(chunk, chunk_idx, file_basename, read_path): results = [] for column in chunk.columns: series = chunk[column] for idx, value in series.items(): try: detect_result = chardet.detect(value) encoding = detect_result['encoding'] if encoding: text = value.decode(encoding) results.append({ "text": text, "encoding": encoding, "row_index": idx, "column_name": column, "chunk_index": chunk_idx }) except (UnicodeDecodeError, LookupError): continue if results: output_filename = f"{file_basename}_chunk{chunk_idx}_results.json" with open(os.path.join(read_path, output_filename), 'w') as f: json.dump(results, f, indent=2) def process_file_with_chunks(file_path, read_path, chunksize=10000): file_basename = os.path.basename(file_path).replace('.parquet', '') # 分块读取Parquet文件 for chunk_idx, chunk in enumerate(pd.read_parquet(file_path, chunksize=chunksize)): process_chunk(chunk, chunk_idx, file_basename, read_path)
4. 减少控制台IO开销
频繁的print操作会因控制台IO阻塞进程,建议将打印内容写入日志文件,或仅打印每个文件的处理统计结果(如成功解码单元格数量)。
额外优化建议
- 直接使用
pyarrow读取Parquet:绕过Pandas的部分封装,直接用pyarrow操作数据,降低性能开销。 - 预分配结果容器:处理前提前创建结果列表,避免动态扩容带来的额外消耗。
- 检查Parquet原始格式:若文件存储时已记录编码信息,可直接读取时指定编码,省去编码检测步骤。
内容的提问来源于stack exchange,提问作者Lenn Box
相关产品推荐
相关产品推荐

