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

Python multiprocessing.pool使用imap_unordered时进程卡住求助

问题描述

我编写了一段使用pool.imap_unordered的多进程Python脚本,在Ubuntu系统运行时出现进程卡住的情况:通过top命令查看发现CPU无负载,进入screen会话后可见进程卡在执行状态。我并未使用存在问题的fork()方法,处理少量数据时不会冻结,但即使是小于5MB的CSV表格数据也会触发该问题。尝试过增加进程数并减少单进程处理数据量,但未解决问题,想请教是否可通过semaphore、lock或其他方式解决,还是需要改用Python并行库?

测试代码1:

import multiprocessing
def my_func(df):
   # modify df here
   # df = df.head(1)
   return df
if __name__ == "__main__":
    df = pd.DataFrame({'a': [2, 2, 1, 1, 3, 3], 'b': [4, 5, 6, 4, 5, 6], 'c': [4, 5, 6, 4, 5, 6]})
    with multiprocessing.Pool(processes = (multiprocessing.cpu_count() - 1)) as pool:
        groups = (g for _, g in df.groupby("a"))
        print(df)
        print(groups)
        out = []
        for res in pool.imap_unordered(my_func, groups):
            out.append(res)
    final_df = pd.concat(out)

测试代码2(改用spawn上下文):

import multiprocessing
def my_func(df):
   # modify df here
   # df = df.head(1)
   return df
if __name__ == "__main__":
    df = pd.DataFrame({'a': [2, 2, 1, 1, 3, 3], 'b': [4, 5, 6, 4, 5, 6], 'c': [4, 5, 6, 4, 5, 6]})
    with multiprocessing.get_context("spawn").Pool(processes = (multiprocessing.cpu_count() - 1)) as pool:
        groups = (g for _, g in df.groupby("a"))
        print(df)
        print(groups)
        out = []
        for res in pool.imap_unordered(my_func, groups):
            out.append(res)
    final_df = pd.concat(out)
排查与解决思路

1. 定位DataFrame序列化问题

multiprocessing传递DataFrame依赖pickle序列化,如果my_func中的操作引入了无法正常pickle的对象(比如自定义lambda、未关闭的文件句柄、第三方库的特殊对象),会导致进程静默卡住。

  • 先将my_func简化为直接返回传入的df,测试是否还卡住。如果恢复正常,逐步加回原有逻辑,定位触发问题的代码段。
  • 手动验证序列化:在主进程中对单个分组的df执行pickle.dumps(g)和pickle.loads(),检查是否抛出异常。

2. 替换imap_unordered为其他任务提交方式

imap_unordered的迭代器式任务提交可能因生成速度与处理速度不匹配导致死锁,可尝试:

  • pool.map():一次性提交所有任务,适合任务量不大的场景
  • pool.apply_async():手动管理任务提交与结果收集,灵活性更高

示例(使用apply_async):

import multiprocessing
import pandas as pd

def my_func(df):
    # 你的处理逻辑
    return df

if __name__ == "__main__":
    df = pd.DataFrame({'a': [2, 2, 1, 1, 3, 3], 'b': [4, 5, 6, 4, 5, 6], 'c': [4, 5, 6, 4, 5, 6]})
    with multiprocessing.get_context("spawn").Pool(processes = multiprocessing.cpu_count() - 1) as pool:
        results = []
        for _, g in df.groupby("a"):
            res = pool.apply_async(my_func, args=(g,))
            results.append(res)
        # 批量收集结果
        out = [res.get() for res in results]
    final_df = pd.concat(out)

3. 用信号量控制并发(针对资源竞争场景)

如果卡住是因为共享资源(如文件、数据库连接)的竞争导致,可通过multiprocessing.Semaphore限制同时访问的进程数:

import multiprocessing
import pandas as pd

# 限制最多4个进程同时执行my_func中的资源密集操作
sem = multiprocessing.Semaphore(4)

def my_func(df):
    with sem:
        # 你的处理逻辑(如文件读写、数据库操作)
        return df

if __name__ == "__main__":
    df = pd.DataFrame({'a': [2, 2, 1, 1, 3, 3], 'b': [4, 5, 6, 4, 5, 6], 'c': [4, 5, 6, 4, 5, 6]})
    with multiprocessing.get_context("spawn").Pool(processes = multiprocessing.cpu_count() - 1) as pool:
        groups = (g for _, g in df.groupby("a"))
        out = list(pool.imap_unordered(my_func, groups))
    final_df = pd.concat(out)

4. 备选并行库方案

如果以上方法无效,可尝试更适配pandas的并行库:

  • swifter:自动根据任务类型选择最优并行方式(单进程/多进程/多线程),无需修改太多原有pandas代码
  • dask:专为大数据集设计的并行计算框架,支持分块处理DataFrame,避免内存过载

5. 系统层面排查

Ubuntu上进程无负载卡住也可能是系统资源限制导致:

  • 用ulimit -n检查文件描述符限制,若处理CSV时打开过多文件,可能触发阻塞
  • 用strace -p <卡住的进程ID>跟踪系统调用,查看进程卡在哪个操作上,定位具体原因

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 02:47:03