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

