pandas read_csv分块后用线程池处理未实现多线程的原因是什么?
问题核心原因
- 线程池不适用于CPU密集型任务:Pandas的底层计算属于CPU密集操作,受Python全局解释器锁(GIL)限制,多线程无法实现真正的并行计算,同一时间只会有一个线程执行Python字节码,所以即使配置了多线程也不会有速度提升。
- 分块读取逻辑本身串行执行:
pd.read_csv返回的分块迭代器,只有在迭代获取下一个块的时候才会实际执行IO读取、解析CSV生成DataFrame的操作。你直接把迭代器传给executor.map,相当于所有块的读取解析都在主线程串行完成,线程池只做了返回dtype的极轻量操作,自然看不到多线程效果。
修复方案
CPU密集型任务需要用进程池绕开GIL,同时把分块读取的逻辑放到子进程中执行,避免主线程串行读块的瓶颈。
修正后代码
import pandas as pd from concurrent.futures import ProcessPoolExecutor def count_file_lines(file_path): """轻量计算文件总行数,不加载全量内容""" with open(file_path, 'r', encoding='latin1') as f: return sum(1 for _ in f) def read_and_check_chunk(args): _input_file_path, _col_index, skiprows, nrows, is_first_chunk = args chunk = pd.read_csv( _input_file_path, na_values=[".", "NA"], skiprows=skiprows, nrows=nrows, encoding="latin1", low_memory=False, sep="\t", usecols=[_col_index], header=0 if is_first_chunk else None ) # 此处可扩展逐行校验dtype的逻辑,当前保留原逻辑返回块的dtype return str(chunk.dtypes.iloc[0]) def read_file_in_chunks(_input_file_path, _col_index, chunksize=10000): # 减去表头行得到数据总行数 total_data_lines = count_file_lines(_input_file_path) - 1 task_args = [] current_offset = 0 while current_offset < total_data_lines: is_first = current_offset == 0 # 首块跳过0行,包含表头;后续块跳过表头+已读数据行 skiprows = 0 if is_first else (current_offset + 1) task_args.append((_input_file_path, _col_index, skiprows, chunksize, is_first)) current_offset += chunksize _dtypes = [] # 进程池worker数建议和CPU核心数保持一致 with ProcessPoolExecutor(max_workers=8) as executor: for res in executor.map(read_and_check_chunk, task_args): _dtypes.append(res) return _dtypes
注意事项
- 进程池的
max_workers不要超过设备CPU核心数,否则会因为进程切换额外开销降低效率 - 如果需要逐行校验dtype,只需要把校验逻辑写在
read_and_check_chunk函数内即可,所有块的校验操作会并行执行 - 该方案不会一次性加载全量文件内容,内存占用和单分块读取方案一致,同时能利用多核CPU提升处理速度
内容的提问来源于stack exchange,提问作者user1261558
相关产品推荐
相关产品推荐

