Python百万级大数据多进程处理优化及防卡顿问题咨询
问题分析与解决方案
你当前代码的核心坑点
- 全局变量
df在多进程中会被每个进程完整复制一份,百万级数据直接导致内存暴增,这就是电脑卡顿、程序停滞的根本原因 - 线程池处理数据清洗完全无效:Pandas的核心操作都是CPU密集型,Python的GIL锁会让线程无法真正并行
- O(n²)的双重循环逻辑在百万级数据下完全不可行,先优化算法再谈并行才是正道
1. 用Multiprocess处理百万级数据的正确姿势
第一步:优化数据加载与预处理
不要把全量数据一次性塞进内存,采用分块加载+预处理的方式,降低内存压力:
import pandas as pd import multiprocessing as mp import time # 单块数据预处理函数 def preprocess_chunk(chunk): return chunk.dropna().reset_index(drop=True) # 分块加载并清洗数据 def load_and_preprocess(file_path, chunksize=10000): cleaned_chunks = [] for chunk in pd.read_csv(file_path, skipinitialspace=True, encoding='utf8', engine='python', chunksize=chunksize): cleaned_chunks.append(preprocess_chunk(chunk)) return pd.concat(cleaned_chunks, ignore_index=True) # 加载预处理后的数据 df = load_and_preprocess('rows.csv')
第二步:重构并行比较逻辑
先明确你的比较目标:如果是找重复项,直接用Pandas内置的duplicated()比两两循环快几个数量级;如果必须做两两比较,要把任务拆成互不重叠的子块,避免重复计算:
# 拆分比较任务:按批次生成需要比较的索引对 def split_compare_tasks(total_rows, batch_size=100): task_batches = [] for i in range(total_rows): # 只比较i之后的行,避免重复计算 task_batches.append((i, range(i+1, total_rows))) # 按批次拆分,防止单个任务过大 if len(task_batches) % batch_size == 0: yield task_batches task_batches = [] if task_batches: yield task_batches # 单进程执行的比较任务 def compare_batch(task_batch): results = [] for i, t_range in task_batch: product_i = df.iloc[i]['Product'] for t in t_range: try: product_t = df.iloc[t]['Product'] # 替换成你的实际比较逻辑(比如相似度计算、相等判断) if product_i == product_t: results.append((i, t, product_i)) except Exception as e: results.append(f"Error comparing {i} & {t}: {str(e)}") return results # 并行执行比较流程 def run_parallel_compare(df, max_workers=None): start_time = time.time() total_rows = len(df) task_batches = split_compare_tasks(total_rows) # 动态设置进程数,默认用CPU核心数 if max_workers is None: max_workers = mp.cpu_count() with mp.Pool(processes=max_workers) as pool: all_results = pool.map(compare_batch, task_batches) # 合并所有结果 final_results = [] for batch_results in all_results: final_results.extend(batch_results) print(f"总耗时: {time.time() - start_time:.2f}秒") return final_results
第三步:避免多进程内存爆炸
- 不要让每个进程持有全量df:如果内存不足以支撑多进程复制全量数据,可以把df导出为Parquet格式(比CSV更高效),每个进程按需读取指定行
- 用共享内存(可选):Python 3.8+支持
pandas的共享内存特性,能让多进程共享同一份数据内存,但实现较复杂,优先用分块处理方案
2. 如何避免电脑卡顿
- 限制进程数:不要盲目开20个进程,按CPU核心数来设置(后面会讲具体数值)
- 合理设置分块大小:百万级数据建议chunksize设为10000-50000,太小会频繁IO,太大则内存占用过高
- 删掉无意义的打印:原代码里的
print会严重拖慢并行速度,换成日志或最后统一输出结果 - 监控系统资源:用任务管理器(Windows)或
top(Linux/macOS)看CPU、内存占用,内存超过80%就调小分块大小或进程数 - 优先优化算法:O(n²)的逻辑在百万级数据下即使并行也会慢到离谱,能用Pandas内置函数或向量化操作解决的,绝不用循环
3. 可设置的最大进程数是多少
- CPU密集型任务(比如你的数据比较):进程数设为CPU核心数或核心数+1,比如8核CPU就开8-9个进程,超过核心数会导致进程频繁切换,反而变慢
- 内存限制优先:如果每个进程需要占用1G内存,总内存16G的话,最多开10个左右(留6G给系统),实际最大进程数要结合内存承受能力调整
- 动态获取核心数:用
mp.cpu_count()可以直接获取当前CPU核心数,代码里可以动态设置:max_workers = mp.cpu_count()
最终优化后的完整代码
import pandas as pd import multiprocessing as mp import time def preprocess_chunk(chunk): return chunk.dropna().reset_index(drop=True) def load_and_preprocess(file_path, chunksize=10000): cleaned_chunks = [] for chunk in pd.read_csv(file_path, skipinitialspace=True, encoding='utf8', engine='python', chunksize=chunksize): cleaned_chunks.append(preprocess_chunk(chunk)) return pd.concat(cleaned_chunks, ignore_index=True) def split_compare_tasks(total_rows, batch_size=100): task_batches = [] for i in range(total_rows): task_batches.append((i, range(i+1, total_rows))) if len(task_batches) % batch_size == 0: yield task_batches task_batches = [] if task_batches: yield task_batches def compare_batch(task_batch): global df results = [] for i, t_range in task_batch: product_i = df.iloc[i]['Product'] for t in t_range: try: product_t = df.iloc[t]['Product'] # 替换成你的实际比较逻辑 if product_i == product_t: results.append((i, t, product_i)) except Exception as e: results.append(f"Error: {i}-{t} -> {str(e)}") return results def run_parallel_compare(df, max_workers=None): if max_workers is None: max_workers = mp.cpu_count() start_time = time.time() total_rows = len(df) task_batches = split_compare_tasks(total_rows) with mp.Pool(processes=max_workers) as pool: all_results = pool.map(compare_batch, task_batches) final_results = [] for batch_results in all_results: final_results.extend(batch_results) print(f"总耗时: {time.time() - start_time:.2f}秒") return final_results if __name__ == '__main__': # Windows系统必须加这个判断,否则多进程会重复执行代码 df = load_and_preprocess('rows.csv') results = run_parallel_compare(df) # 将结果保存到文件 pd.DataFrame(results, columns=['Index1', 'Index2', 'Product']).to_csv('compare_results.csv', index=False)
内容的提问来源于stack exchange,提问作者pln
相关产品推荐
相关产品推荐

