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

使用Python concurrent.futures处理大量输入时进程停滞问题排查

问题根源分析

1. 初始代码的问题

当输入规模达到100万时,executor.map()会尝试一次性处理整个迭代器。在Windows系统(Python 3.10默认用spawn方式创建子进程)下,每个子进程启动时都需要重新导入主模块,同时一次性提交百万级任务会导致:

  • 大量进程同时启动,耗尽系统进程资源或内存
  • 任务参数的序列化/反序列化开销急剧增大,拖慢启动速度甚至直接崩溃

2. 批量拆分代码的问题

你每次循环都创建新的ProcessPoolExecutor,这会引发两个关键问题:

  • 频繁创建和销毁进程池,带来巨大的系统开销(进程启动、销毁本身就属于高耗时操作)
  • 在Windows的spawn模式下,每次创建进程池都会重复导入主模块,累积的延迟会让程序看起来卡住甚至无法启动
  • 最关键的是你的代码没有if __name__ == '__main__':保护,子进程启动时会重新执行整个脚本,陷入无限创建进程的死循环,直接导致程序卡死
修复方案

核心修复:添加主程序入口保护

在Windows系统下使用多进程,必须把主逻辑放在if __name__ == '__main__':块中,否则子进程会重复执行脚本代码,引发致命异常。

优化方案:复用单个进程池

不要每次批量任务都重建进程池,而是创建一个进程池后分批次提交任务,彻底减少进程创建的开销。

优化后的代码示例:

import concurrent.futures
import math

def some_math(x):
    y = 3*x**2 + 5*x + 7
    return math.log(y)

if __name__ == '__main__':
    up_to = 1_000_000
    batch_size = 20_000
    results = []
    
    # 复用单个进程池,指定进程数(建议等于CPU核心数)
    with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
        for idx in range(0, up_to, batch_size):
            low = idx
            high = min(low + batch_size, up_to)
            inputs = range(low, high)
            # 提交当前批次任务并收集结果
            batch_results = list(executor.map(some_math, inputs))
            results.extend(batch_results)

额外优化建议

  • 控制进程池大小:max_workers不要设置过大,CPU密集型任务建议设为CPU核心数(可通过os.cpu_count()获取),避免进程过多导致上下文切换开销暴增
  • 使用chunksize参数:对于executor.map(),可以指定chunksize参数,让每个子进程一次性处理多个任务,减少进程间通信的开销。比如4个进程处理2万条任务时,可设置chunksize=5000:
    batch_results = list(executor.map(some_math, inputs, chunksize=batch_size//4))
    
  • 避免一次性加载所有数据:如果处理的是文件而非数值范围,确保每次只加载当前批次的文件路径,不要一次性把所有路径加载到内存

内容的提问来源于stack exchange,提问作者meliksahturker

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:01:14