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

Python多线程.parquet单元格编码检测脚本CPU利用率低的优化问题

问题根源分析
  1. 任务粒度太小:每个单元格作为独立进程任务,进程间数据传递(IPC)的开销远大于编码检测的计算开销,导致CPU大部分时间消耗在数据传输上,无法有效利用。
  2. 单文件串行处理:脚本逐个读取Parquet文件、生成任务列表后再提交到进程池,进程池在等待文件读取、任务列表生成时处于闲置状态。
  3. 低效的文件IO:每个成功解码的单元格都创建独立JSON文件,还需循环判断文件名是否存在,大量小文件IO会严重阻塞进程,拖慢整体速度。
  4. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 01:10:56