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

Python中pool.starmap()运行远慢于pool.map()的原因排查

问题分析与解决方案

核心变慢原因

  • 并行利用完全失效:原pool.map会把大列表拆成与CPU核心数匹配的多个chunk,所有核心同时并行处理;而你的starmap代码每次只传递1个完整任务(比如整个list_a1+list_b1),进程池只能启动1个进程处理,且四次starmap调用是串行执行,相当于单进程跑完全部任务,完全浪费了多核优势,这是速度暴跌的主要原因。
  • chunksize参数无效:chunksize是当任务列表有多个任务时,指定每个进程处理的任务数。但你每次的任务列表只有1个元素,设置chunksize=multiprocessing.cpu_count()没有任何作用,无法拆分任务到多个进程。
  • 查找效率极低:item in list_B是O(n)时间复杂度的操作,500万条数据下会产生大量重复遍历,这是隐性的性能瓶颈。

针对性解决方案

1. 优先优化查找效率(最关键)

把list_B转换成集合,利用集合O(1)的查找特性,将整体时间复杂度从O(n*m)降至O(n):

def listCompare(list_A, set_B):
    # 用列表推导式简化代码,执行效率更高
    return [item for item in list_A if item in set_B]

2. 重构任务拆分,充分利用多核

沿用原chunks函数拆分大列表,将每个chunk与对应的集合化list_B组成任务元组,让starmap可以并行处理多个任务:

import multiprocessing

def chunks(lst, n):
    """将列表拆分为n个均匀的chunk"""
    for i in range(n):
        yield lst[i::n]

def listCompare(list_A, set_B):
    return [item for item in list_A if item in set_B]

if __name__ == "__main__":
    # 替换为你的实际数据
    list_a1 = [...]
    list_b1 = [...]
    
    # 转换为集合提升查找效率
    set_b1 = set(list_b1)
    # 拆分list_a1为与CPU核心数匹配的chunk
    chnks_a1 = chunks(list_a1, multiprocessing.cpu_count())
    # 生成任务列表:每个chunk对应同一个set_b1
    tasks_a1 = [(chunk, set_b1) for chunk in chnks_a1]
    
    with multiprocessing.Pool() as pool:
        # 一次性并行处理所有chunk
        results_a1 = pool.starmap(listCompare, tasks_a1)
        output_list1 = []
        for res in results_a1:
            output_list1.extend(res)
    
    # 同理处理list_a2/list_b2、list_a3/list_b3等组
    # ...

3. 可选:批量处理多组任务

如果四组数据可以同时处理,还可以把所有组的任务合并成一个大任务列表,一次starmap调用完成,减少进程池的启动/关闭开销:

if __name__ == "__main__":
    # 定义所有待处理的组
    data_groups = [
        (list_a1, list_b1),
        (list_a2, list_b2),
        (list_a3, list_b3),
        (list_a4, list_b4),
    ]
    
    all_tasks = []
    cpu_count = multiprocessing.cpu_count()
    
    for list_a, list_b in data_groups:
        set_b = set(list_b)
        chnks_a = chunks(list_a, cpu_count)
        all_tasks.extend([(chunk, set_b) for chunk in chnks_a])
    
    with multiprocessing.Pool() as pool:
        all_results = pool.starmap(listCompare, all_tasks)
        
        # 按组拆分结果(需要提前记录每组的chunk数量,这里简化演示)
        output_list1 = []
        output_list2 = []
        # ... 自行处理结果分组

补充说明

  • 必须加上if __name__ == "__main__":保护,避免Windows系统下的多进程递归启动问题。
  • 集合转换不会影响元素的频率保留:我们只是用集合快速判断元素是否存在,最终结果还是保留原list_A中的元素顺序和出现次数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 08:36:19