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

Python多进程池循环处理数据挂起问题及终止方法咨询

解决进程池循环使用时的挂起问题

首先咱们来拆解你代码里的核心问题:你在with语句中手动调用了pool.close()和pool.terminate(),这完全是画蛇添足,甚至是导致挂起的关键原因。

为什么会挂起?

with Pool() as pool这个语法本身就会帮你自动管理进程池的生命周期:当代码块执行完毕时,它会自动调用pool.close()(禁止进程池接受新任务)和pool.join()(等待所有子进程完成任务),最后优雅地销毁进程池。

而你手动加的pool.terminate()是强制杀死所有子进程,不管它们有没有完成任务。这会导致部分子进程的资源无法正常释放,随着循环次数增加,未释放的资源越积越多,最终导致程序挂起。缩小数据块大小只是延迟了资源耗尽的时间,本质问题没解决。

正确的两种处理方式

方式1:复用进程池(推荐,更高效)

没必要每次循环都新建进程池,直接在循环外创建一次,所有分块任务都用这同一个进程池处理,既减少进程创建销毁的开销,也避免资源泄漏:

import pandas as pd
import numpy as np
from multiprocess import Pool

df = pd.read_csv('paths.csv')

def do_something(user):
    v = df[df['userId'] == user]
    return v

if __name__ == '__main__':
    users = df['userId'].unique()
    n_chunks = round(len(users)/40)
    subsets = [users[i:i+n_chunks] for i in range(0, len(users), n_chunks)]
    chunk_counter = 0
    
    # 循环外创建进程池,所有任务复用
    with Pool() as pool:
        for user_subset in subsets:
            chunk_counter += 1
            print(f'Beginning to process chunk {chunk_counter}...')
            frames = pool.map(do_something, user_subset)
            print(f'Completed processing chunk {chunk_counter}.')
    # with代码块结束后,自动关闭并回收进程池

方式2:每次循环重建进程池(仅当有特殊需求时用)

如果你确实需要每次循环都新建进程池,那要抛弃terminate(),改用close()+join()来确保子进程完成任务后再正常销毁:

import pandas as pd
import numpy as np
from multiprocess import Pool

df = pd.read_csv('paths.csv')

def do_something(user):
    v = df[df['userId'] == user]
    return v

if __name__ == '__main__':
    users = df['userId'].unique()
    n_chunks = round(len(users)/40)
    subsets = [users[i:i+n_chunks] for i in range(0, len(users), n_chunks)]
    chunk_counter = 0
    
    for user_subset in subsets:
        chunk_counter += 1
        print(f'Beginning to process chunk {chunk_counter}...')
        pool = Pool()
        frames = pool.map(do_something, user_subset)
        # 先关闭进程池,禁止新任务进入
        pool.close()
        # 等待所有子进程完成当前任务
        pool.join()
        # join完成后,进程池会自动销毁,资源正常释放
        print(f'Completed processing chunk {chunk_counter}.')

额外优化建议

你的do_something函数引用了全局的df,在多进程模式下,每个子进程都会复制一份df到自己的内存空间,如果df很大,会导致内存占用飙升。可以考虑:

  • 把df的必要列转换成共享内存对象(比如用multiprocessing.Array或pandas的共享内存工具)
  • 或者在do_something中只传递需要的用户ID和数据片段,避免全局变量复制

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:00:36