Python多进程线程数增加时性能下降问题排查求助
多进程并行性能随线程数增加下降的排查与优化问题
我拥有一台24核、每核2线程的机器,正尝试优化以下代码以实现并行执行,但发现线程数在7-8个以内时扩展性良好,超过该数量后性能持续恶化。通过性能分析发现,每个线程执行相同代码的耗时随线程数增加而显著上升。
待优化代码
import argparse import glob import h5py import numpy as np import pandas as pd import xarray as xr from tqdm import tqdm import time import datetime from multiprocessing import Pool, cpu_count, Lock import multiprocessing import cProfile, pstats, io def process_parcel_file(f, bands, mask): start_time = time.time() test = xr.open_dataset(f) print(f"Elapsed in process_parcel_file for reading dataset: {time.time() - start_time}") start_time = time.time() subset = test[bands + ['SCL']].copy() subset = subset.where(subset != 0, np.nan) if mask: subset = subset.where((subset.SCL >= 3) & (subset.SCL < 7)) subset = subset[bands] # Adding a new dimension week_year and performing grouping subset['week_year'] = subset.time.dt.strftime('%Y-%U') subset = subset.groupby('week_year').mean().sortby('week_year') subset['id'] = test['id'].copy() # Store the dates and counting pixels for each parcel dates = subset.week_year.values n_pixels = test[['id', 'SCL']].groupby('id').count()['SCL'][:, 0].values.reshape(-1, 1) # Converting to dataframe grouped_sum = subset.groupby('id').sum() ids = grouped_sum.id.values grouped_sum = grouped_sum.to_array().values grouped_sum = np.swapaxes(grouped_sum, 0, 1) grouped_sum = grouped_sum.reshape((grouped_sum.shape[0], -1)) colnames = ["{}_{}".format(b, str(x).split('T')[0]) for b in bands for x in dates] + ['count'] values = np.hstack((grouped_sum, n_pixels)) df = pd.DataFrame(values, columns=colnames) df.insert(0, 'id', ids) print(f"Elapsed in process_parcel_file til end: {time.time() - start_time}") return df def fs_creation(input_dir, out_file, labels_to_keep=None, th=0.1, n=64, days=5, total_days=180, mask=False, mode='s2', method='patch', bands=['B02', 'B03', 'B04', 'B05', 'B06', 'B07', 'B08', 'B8A', 'B11', 'B12']): files = glob.glob(input_dir) times_pool = [] # For storing execution times times_seq = [] cpu_counts = list(range(2, multiprocessing.cpu_count() + 1, 4)) # The different CPU counts to use for count in cpu_counts: print(f"Executing with {count} threads") if method == 'parcel': start_pool = time.time() with Pool(count) as pool: arguments = [(f, bands, mask) for f in files] dfs = list(tqdm(pool.starmap(process_parcel_file, arguments), total=len(arguments))) end_pool = time.time() start_seq = time.time() dfs = pd.concat(dfs) dfs = dfs.groupby('id').sum() counts = dfs['count'].copy() dfs = dfs.div(dfs['count'], axis=0) dfs['count'] = counts dfs.drop(index=-1).to_csv(out_file) end_seq = time.time() times_pool.append(end_pool - start_pool) times_seq.append(end_seq - start_seq) pd.DataFrame({'CPU_count': cpu_counts, 'Time pool': times_pool, 'Time seq' : times_seq}).to_csv('cpu_times.csv', index=False) return 0
性能对比示例
2线程执行耗时
Elapsed in process_parcel_file for reading dataset: 0.012271404266357422 Elapsed in process_parcel_file til end: 1.6681673526763916 Elapsed in process_parcel_file for reading dataset: 0.014229536056518555 Elapsed in process_parcel_file til end: 1.5836331844329834
22线程执行耗时
Elapsed in process_parcel_file for reading dataset: 0.17968058586120605 Elapsed in process_parcel_file til end: 12.049026727676392 Elapsed in process_parcel_file for reading dataset: 0.052398681640625 Elapsed in process_parcel_file til end: 6.014119625091553
补充信息
- 单份文件约80MB,共451份;
- 已使用memory_profiler与psutil对代码各环节进行性能分析,有相关CSV结果及线程工作情况报告;
- 无法提供最小可复现示例。
恳请提供指导建议以定位问题根源并优化代码性能。
内容的提问来源于stack exchange,提问作者Norhther
相关产品推荐
相关产品推荐

