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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 05:25:14