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

如何用Python/Dask/多进程并行化处理分子DataFrame分组矩阵计算?

优化分子矩阵计算的并行化方案

现有多进程代码的问题

你的multiprocessing实现效率低下的核心原因是每个ID循环内都创建新的进程池,进程的创建与销毁开销远超过单个ID内组合计算的耗时。此外,每个ID的Pose数量有限(实际约30个),并行计算少量任务的收益完全被进程管理成本抵消。

优化方案

1. 正确的多进程实现:按ID级并行

将每个ID的完整矩阵计算作为独立任务,在全局创建一次进程池,批量处理所有ID。这样进程可以复用,避免重复创建的开销。

示例代码:

import pandas as pd
import itertools
import multiprocessing
import numpy as np
from tqdm import tqdm

# 生成示例数据(保留原逻辑)
df = pd.DataFrame(columns=['ID', 'Pose ID'])
ids = ['ID' + str(i) for i in range(1, 6)]
pose_ids = ['Pose ' + str(i) for i in range(1, 11)]
df_list = []
for i in ids:
    temp_df = pd.DataFrame({'ID': [i] * 10, 'Pose ID': pose_ids})
    df_list.append(temp_df)
df = pd.concat(df_list)

# 预处理:按ID分组,提前准备每个组的数据
grouped = df.groupby('ID')
id_groups = [(id_val, group) for id_val, group in grouped]

# 单个ID的矩阵计算函数
def process_single_id(args):
    id_val, group = args
    poses = group['Pose ID'].tolist()
    n = len(poses)
    # 预处理Pose到索引的映射,避免重复查找
    pose_to_idx = {pose: idx for idx, pose in enumerate(poses)}
    # 用NumPy数组初始化矩阵,比Pandas DataFrame更快
    matrix = np.zeros((n, n), dtype=object)  # 根据实际计算结果类型调整dtype
    
    # 计算所有组合
    for pose1, pose2 in itertools.combinations(poses, 2):
        idx1 = pose_to_idx[pose1]
        idx2 = pose_to_idx[pose2]
        # 替换为实际的分子计算逻辑
        result = pose1 + pose2
        matrix[idx1, idx2] = result
        matrix[idx2, idx1] = result
    
    # 转成Pandas DataFrame(如果需要)
    return (id_val, pd.DataFrame(matrix, index=poses, columns=poses))

# 并行处理所有ID
if __name__ == '__main__':
    # 根据CPU核心数设置进程数
    with multiprocessing.Pool(processes=multiprocessing.cpu_count()) as pool:
        # 用tqdm显示进度
        results = list(tqdm(pool.imap_unordered(process_single_id, id_groups), total=len(id_groups)))
    
    # 整理成字典
    calculated_dfs = {id_val: mat for id_val, mat in results}

2. Dask解决方案:确保ID分组在同一分区

Dask的核心问题是让同一ID的所有数据落在同一个分区,你可以通过以下两种方式实现:

方法1:按ID设置索引并重分区

import dask.dataframe as dd
import numpy as np
import itertools

# 将Pandas DataFrame转为Dask DataFrame,按ID设置索引
ddf = dd.from_pandas(df, npartitions=multiprocessing.cpu_count())
ddf = ddf.set_index('ID')

# 定义单个组的计算函数,需指定meta参数
def process_group(group):
    poses = group['Pose ID'].tolist()
    n = len(poses)
    pose_to_idx = {pose: idx for idx, pose in enumerate(poses)}
    matrix = np.zeros((n, n), dtype=object)
    for pose1, pose2 in itertools.combinations(poses, 2):
        idx1 = pose_to_idx[pose1]
        idx2 = pose_to_idx[pose2]
        result = pose1 + pose2
        matrix[idx1, idx2] = result
        matrix[idx2, idx1] = result
    return pd.DataFrame(matrix, index=poses, columns=poses)

# 应用分组计算,meta指定返回类型(示例用空DataFrame)
meta = pd.DataFrame(columns=pose_ids, index=pose_ids)
calculated_dfs = ddf.groupby('ID').apply(process_group, meta=meta).compute()

方法2:使用Dask Bag处理独立分组

将每个ID的组转为Dask Bag的元素,天然保证每个组在同一分区:

from dask.bag import from_sequence
import numpy as np
import itertools

# 将id_groups转为Dask Bag
bag = from_sequence(id_groups, npartitions=multiprocessing.cpu_count())

# 映射处理函数
results_bag = bag.map(process_single_id)

# 计算结果并整理为字典
calculated_dfs = dict(results_bag.compute())

3. 其他性能优化点

  • 减少Pandas索引查找开销:提前创建Pose ID到矩阵索引的映射字典,避免每次用布尔索引查找位置,大幅减少循环内耗时。
  • 用NumPy数组存储中间矩阵:NumPy的数值操作比Pandas更高效,仅在最终结果需要时转为DataFrame。
  • 批量预处理分组数据:提前用df.groupby('ID')准备所有组的数据,避免循环内重复切片DataFrame。
  • 向量化计算逻辑:如果分子相似度计算支持向量化(比如用矩阵运算代替循环),优先使用向量化操作,收益比并行化更显著。

内容的提问来源于stack exchange,提问作者Antoine Lacour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:50:21