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

Python多进程Pool仅首次map执行,后续无响应问题求助

问题排查与修复方案

你遇到的核心问题是索引不匹配,导致后续循环中thread_search函数实际上没有执行任何有效操作,看起来像是map没有运行。另外还有多进程全局变量共享的问题,我逐一给你分析:

1. 索引错误:原数组索引 vs 切片后数组的本地索引

当你处理第2、3组数据时,x是原CSV数组的切片(x = list(csv.reader(...))[lower_bound:upper_bound]),也就是说x的自身索引是从0到len(x)-1,对应原数组的lower_bound到upper_bound位置。

但你生成的divisor_pairs用的是原数组的索引范围:

divisor_pairs = [ [int(lower_bound + (i - 1)*chunk_size) , int(lower_bound + i*chunk_size)] for i in range(1,6143) ]

在thread_search里,你用这个原数组索引去访问x:

current_section = np.array(x[i: i + window_size])

比如第2组的lower_bound = group_step,此时i的起始值是group_step,但x的长度只有group_step,所以x[i:i+window_size]是空数组,整个循环for i in range(divisor_lower, divisor_upper,window_size)没有任何迭代操作,自然不会产生任何输出,看起来就像map没运行。

修复方法:

把divisor_pairs改为使用切片后x的本地索引,去掉lower_bound偏移:

divisor_pairs = [ [int((i - 1)*chunk_size) , int(i*chunk_size)] for i in range(1,6143) ]

2. 全局变量m的多进程共享问题

你用全局变量m做计数,但在多进程中,每个子进程会复制父进程的内存空间,所以每个子进程都有自己的m副本,不是共享的。这就是为什么第一次运行时会重复打印相同的m值(比如多个Encountered m = 1000000)。

修复方法:

使用multiprocessing.Value创建共享计数器,确保所有子进程操作同一个计数:

from multiprocessing import Value

def iterCount(m):
    # 加锁确保原子操作,避免计数混乱
    with m.get_lock():
        m.value += 1
        return m.value

def thread_search(args):
    pair, m = args
    divisor_lower = pair[0]
    divisor_upper = pair[1]
    for i in range(divisor_lower, divisor_upper, window_size):
        current_section = np.array(x[i: i + window_size])
        for row in current_section:
            if (row[2].startswith('NP') ) and checkPep(row[0]):
                shared_list.append(row[[0,1,2]])
        current_m = iterCount(m)
        if not current_m % 1000000:
            print(f'Encountered m = {current_m}', flush = True)

def poolMap(pairs, group, m):
    job = Pool(3)
    print(f'Pool Created')
    print(len(pairs))
    # 将共享计数器和pair打包传递
    job.map(thread_search, [(pair, m) for pair in pairs])
    print('Pool Closed')
    job.close()
    job.join()  # 加入join等待所有子进程结束,确保资源回收

if __name__ == '__main__':
    # 创建共享整数计数器,初始值0
    m = Value('i', 0)
    for group in [1,2,3]:
        x = None
        lower_bound = int((group - 1)*group_step)
        upper_bound = int(group*group_step)
        x = list(csv.reader(open(pa_table_name,"rt", encoding = "utf-8"), delimiter = "\t"))[lower_bound:upper_bound]
        print(len(x))
        divisor_pairs = [ [int((i - 1)*chunk_size) , int(i*chunk_size)] for i in range(1,6143) ]
        poolMap(divisor_pairs, group, m)

3. 额外建议:Pool的资源回收

虽然map是阻塞式调用,但在close()后加上join()可以确保所有子进程完全结束,避免资源泄漏,这也是上面修复代码里加入job.join()的原因。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 06:29:01