Python 多处理中线程增加导致性能下降

Nor*_*her 5 python parallel-processing multithreading multiprocessing python-multiprocessing

我有一台 24 个核心的机器,每个核心有 2 个线程。我正在尝试优化以下代码以实现并行执行。但是,我注意到代码的性能在达到一定数量的线程后开始下降。

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
Run Code Online (Sandbox Code Playgroud)

执行代码时,它可以很好地扩展到 7-8 个线程左右,但之后,性能开始恶化。我已经分析了代码,似乎每个线程都需要更多时间来执行相同的代码。

例如,有 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
Run Code Online (Sandbox Code Playgroud)

然而,对于 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
Run Code Online (Sandbox Code Playgroud)

我很难理解为什么线程越多性能就会下降。我已经验证系统具有所需的核心和线程数量。

我将不胜感激任何指导或建议来帮助我确定此问题的原因并优化代码以获得更好的性能。

对我来说,提供一个最小的工作示例确实很难,所以请考虑到这一点。

先感谢您。

编辑:每个文件大约 80MB。我有 451 个文件。我添加了以下代码来分析该函数:

...
    start_time = time.time()
    mem_usage_start = memory_usage(-1, interval=0.1, timeout=None)[0]
    cpu_usage_start = psutil.cpu_percent(interval=None)
    test = xr.open_dataset(f)
    times['read_dataset'] = time.time() - start_time
    memory['read_dataset'] = memory_usage(-1, interval=0.1, timeout=None)[0] - mem_usage_start
    cpu_usage['read_dataset'] = psutil.cpu_percent(interval=None) - cpu_usage_start
...
Run Code Online (Sandbox Code Playgroud)

以类似的方式为每行添加更多代码。我使用了库memory_profiler和psutil,并且我有每个线程的信息。包含结果的 CSV 可在此处获取: https://wetransfer.com/downloads/44df14ea831da7693300a29d8e0d4e7a20230703173536/da04a0 结果用所选的 cpu 数量标识函数中的每一行,因此每一行都是一个线程。

编辑2:

这里我有一个数据子集的报告,您可以在其中清楚地看到每个线程正在做什么,以及某些线程如何比其他线程获得更少的工作:

https://wetransfer.com/downloads/259b4e42aae6dd9cda5a22d576aba29520230717135248/ae3f88

mic*_*ses 0

长话短说

\n

替换xarray.open_dataset为xarray.load_dataset立即将数据加载到内存中;前者非常懒惰,而后者则有点急切(又名贪婪,严格)。
\n[如果这个答案被接受]

\n

解释:

\n

[冒着重复之前评论/答案的风险,]这似乎是 I/O 吞吐量限制,太多线程/进程尝试同时从同一设备读取数据,具有讽刺意味的是,这会减慢所有同级线程/进程的速度。

\n

至于造成这种情况的原因,我认为是第 17 行对xarray\'s 的调用,通过对文档的某些特定解释,它可能看起来是在延迟加载数据:open_datasettest = xr.open_dataset(f)

\n
\n

数据始终从 netCDF 文件延迟加载。您可以对 Dataset 和 DataArray 对象进行操作、切片和子集化,并且在您尝试执行某种实际计算之前,不会将任何数组值加载到内存中。

\n
\n

如果这是真的,它可以解释您所表现出的症状 - 451 个文件最初由“空”对象表示(这花费的时间几乎可以忽略不计),并且随着这些对象的数据被读取(或操作) )几乎同时 - 存储设备受到数万(或数百万,取决于块大小)的读取请求的攻击。

\n

有一个提示,共四段:

\n
\n

Xarray\xe2\x80\x99s 远程或磁盘数据集的延迟加载通常是可取的,但并不总是可取的。在执行计算密集型操作之前,通常最好通过调用 Dataset.load() 方法将 Dataset(或 DataArray)完全加载到内存中。

\n
\n

或者,尝试使用load_dataset。

\n