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

如何在含pool.imap_unordered的Python多进程脚本中正确使用Lock()

在multiprocessing.Pool.imap_unordered中正确使用Lock的方法

首先明确:只有当你的my_func中存在多进程共享资源访问(比如写入同一个文件、修改共享内存变量、避免控制台输出混乱)时,才需要使用锁。如果只是单独处理每个DataFrame分组,不需要加锁。

下面是针对你的脚本的两种正确实现方式:

方式一:使用Manager传递共享锁

通过multiprocessing.Manager()创建可跨进程共享的锁,将锁与每个分组数据打包成元组传入任务:

import multiprocessing
import pandas as pd

def my_func(args):
    df, lock = args
    
    # 临界区操作:需要加锁保护的代码块
    with lock:
        print(f"正在处理分组 a={df['a'].iloc[0]}")
    
    # 处理DataFrame逻辑
    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.Manager() as manager:
        lock = manager.Lock()
        # 将每个分组和锁打包为参数
        groups_with_lock = ((g, lock) for _, g in df.groupby("a"))
        
        with multiprocessing.Pool() as pool:
            out = []
            for res in pool.imap_unordered(my_func, groups_with_lock):
                out.append(res)
    
    final_df = pd.concat(out)
    print(final_df)

方式二:通过Pool初始化器传递锁

利用Pool的initializer和initargs参数,在每个子进程启动时注入锁,避免每次任务都传递锁参数:

import multiprocessing
import pandas as pd

# 全局变量用于存储子进程中的锁引用
lock = None

def init_worker(shared_lock):
    global lock
    lock = shared_lock

def my_func(df):
    # 使用全局锁保护临界区
    with lock:
        print(f"正在处理分组 a={df['a'].iloc[0]}")
    
    # 处理DataFrame逻辑
    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]})
    
    lock = multiprocessing.Lock()
    # 创建进程池时传入初始化函数和锁参数
    with multiprocessing.Pool(initializer=init_worker, initargs=(lock,)) as pool:
        groups = (g for _, g in df.groupby("a"))
        out = []
        for res in pool.imap_unordered(my_func, groups):
            out.append(res)
    
    final_df = pd.concat(out)
    print(final_df)

关键注意点

  • 不能直接传递普通的multiprocessing.Lock()给进程池任务:子进程会复制父进程内存,锁的状态无法跨进程同步,必须用Manager.Lock()或初始化器方式传递共享锁实例。
  • with lock:上下文管理器会自动处理锁的获取和释放,避免手动操作导致死锁。
  • imap_unordered返回结果的顺序是任务完成顺序,锁的使用不影响结果拼接,最终pd.concat(out)仍能得到正确的合并结果。

内容的提问来源于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 13:47:39