Как правильно использовать multiprocessing для при работе с данными?

Возникла у меня проблема. Есть большие мощности, но по умолчанию, только для обучения нейронных сетей, мы используем все ядра. У меня есть функции связанные с препроцессингом данных и другие функции, которые я бы хотел разделить между всеми ядрами, для более быстрого их выполнения и потом соединить результат в одном массиве.

Вот моя функция Вот как я это делаю сейчас:

tresh_list = []
for i in range (1,41):
# def multiproc_func(X, y):
    print('----------------------------------------')
    print('ITERATION: {}'.format(i))
    (X_train_t, y_train_t, X_train_v, y_train_v, X_test, y_test) = data_split(X, y)
    model = model_create(X)
    print('-----------------------------------------')
    print('Train the model')
    model.fit(X_train_t, y_train_t, batch_size=1024, epochs=5, validation_data = (X_train_v, y_train_v))
    print('Training finished')
    print('-----------------------------------------')
    print('Metrics for train set, treshold=0.5')
    print('-----------------------------------------')
    t_pred_t, pred_t = predict(X_train_t, y_train_t)
    print('-----------------------------------------')
    print('Metrics for validation set, treshold=0.5')
    print('-----------------------------------------')
    y_pred_v, pred_v = predict(X_train_v, y_train_v)
    print('-----------------------------------------')
    print('Metrics for test set without, treshold=0.5')
    print('-----------------------------------------')
    y_pred_test, pred_test = predict(X_test, y_test)
    print('-----------------------------------------')
    print('Started threshold-precision distribution. It can get some time...')
    (prec_train, _) = tresh_prec_distribution(y_train_t, pred_t)
    (prec_valid, _) = tresh_prec_distribution(y_train_v, pred_v)
    print('Treshold-precision distrubution finished')
    prec_valid_train = [x + y for x, y in zip(prec_train, prec_valid)]
    prec_id = argmax(prec_valid_train)
    # Drop boundary values
    prec_valid_train.pop(prec_id)
    prec_id = argmax(prec_valid_train)
    prec_valid_train.pop(prec_id)
    best_prec = max(prec_valid_train)
    y_pred = []
    for i in pred_test:
        if i[0] >= treshold[prec_id]:
            y_pred.append(0)
        else:
            y_pred.append(1)
    print('The best treshold: {}'.format(treshold[prec_id]))
    print('Precision score: {}'.format(precision_score(y_test, y_pred)))
    print('Recall score: {}'.format(recall_score(y_test, y_pred)))
    tresh_list.append(treshold[prec_id])

Логика такая: на каждой итерации я определяю нужное мне значение и потом его добавляю в список tresh_list. Но, как я говорил все операции выполняются только на одном ядре, по этому процесс занимает часы.

Как решить эту задачу с использованием мультипроцессинга? Т. е. мне надо снача распарралелить процессы, получить значение treshold[prec_id] для каждого из процесса, и собрать все значения treshold[prec_id] в одном списке tresh_list Я уже читал документацию и смотрел примеры, но не могу понять как их применить в текущей ситуации. Спасибо


Ответы (1 шт):

Автор решения: n1tr0xs

Вы можете использовать этот код для распараллеливания:

from threading import Thread
tresh_list = []
def multiproc_func(X, y):
    global tresh_list
    # Ваш код

todo=[]
for i in range(1,41): # создаем список процессов
    todo.append( Thread(target= multiproc_func, args=(X, y)) )
for th in todo: # стартуем каждый процесс
    th.start()
for th in todo: # запускаем каждый процесс
    th.join()
→ Ссылка