# Параллельные вычисления

__Автор задач: Блохин Н.В. (NVBlokhin@fa.ru)__

Материалы:
* Макрушин С.В. Лекция "Параллельные вычисления"
* https://docs.python.org/3/library/multiprocessing.html
    * https://docs.python.org/3/library/multiprocessing.html#multiprocessing.Process
    * https://docs.python.org/3/library/multiprocessing.html#multiprocessing.pool.Pool
    * https://docs.python.org/3/library/multiprocessing.html#multiprocessing.Queue
* https://pandas.pydata.org/docs/reference/api/pandas.read_csv.html
* https://numpy.org/doc/stable/reference/generated/numpy.array_split.html
* https://nalepae.github.io/pandarallel/
    * https://github.com/nalepae/pandarallel/blob/master/docs/examples_windows.ipynb
    * https://github.com/nalepae/pandarallel/blob/master/docs/examples_mac_linux.ipynb

## Задачи для совместного разбора

In [24]:
# !pip install pandarallel

1. Посчитайте, сколько раз встречается буква "a" в файлах ["xaa", "xab", "xac", "xad"]. 

In [2]:
import multiprocessing

files = [f"{name}.txt" for name in ["xaa", "xab", "xac", "xad"]]

2. Выведите на экран слова из файла words_alpha, в которых есть две или более буквы "e" подряд.

In [28]:
import pandas as pd

words = (
    pd.read_csv("10_multiprocessing_data/words_alpha.txt", header=None)[0]
    .dropna()
    .sample(frac=1, replace=True)
)

## Лабораторная работа 10

__При решении данных задач не подразумевается использования циклов или генераторов Python в ходе работы с пакетами `numpy` и `pandas`, если в задании не сказано обратного. Решения задач, в которых для обработки массивов `numpy` или структур `pandas` используются явные циклы (без согласования с преподавателем), могут быть признаны некорректными и не засчитаны.__

In [1]:
import pandas as pd
import numpy as np
import multiprocessing as mp

1\. В каждой строке файла `tag_nsteps.csv` хранится информация о тэге рецепта и количестве шагов в этом рецепте в следующем виде:

```
tags,n_steps
hungarian,2
european,6
occasion,4
pumpkin,4
................
```

Всего в исходном файле хранится чуть меньше, чем 71 млн, строк. Разбейте файл `tag_nsteps.csv` на несколько (например, 8) примерно одинаковых по объему файлов c названиями `tag_nsteps_*.csv`, где вместо символа `*` указан номер очередного файла. Каждый файл имеет структуру, аналогичную оригинальному файлу (включая заголовок).

__Важно__: здесь и далее вы не можете загружать в память весь исходный файл сразу. 

In [7]:
n_files = 10

In [2]:
files = []
for i in range(n_files):
    files.append(open(f"tag_nsteps_{i}.csv", "w"))
    files[i].write("tags,n_steps\n")
with open("tag_nsteps.csv", "r") as f_main:
    line = f_main.readline()
    line = f_main.readline()
    i = 0
    while line:
        files[i].write(line)
        i = (i + 1) % n_files
        line = f_main.readline()
for i in range(n_files):
    files[i].close()

2\. Напишите функцию, которая принимает на вход название файла, созданного в результате решения задачи 1, считает для каждого тэга сумму по столбцу `n_steps` и количество строк c этим тэгом, и возвращает результат в виде словаря. Ожидаемый вид итогового словаря:

```
{
    '1-day-or-more': {'sum': 56616, 'count': 12752},
    '15-minutes-or-less': {'sum': 195413, 'count': 38898},
    '3-steps-or-less': {'sum': 187938, 'count': 39711},
    ....
}
```

Примените данную функцию к каждому файлу, полученному в задании 1, и соберите результат в виде списка словарей. Не используйте параллельных вычислений. 

Выведите на экран значение по ключу "30-minutes-or-less" для каждого из словарей.

In [8]:
def get_tag_sum_count_from_file(file: str) -> dict:
    df = pd.read_csv(file)
    df2 = df.groupby("tags")["n_steps"].agg(sum = 'sum', count = 'count')
    return df2.to_dict('index')

In [9]:
res2 = [get_tag_sum_count_from_file(f"tag_nsteps_{i}.csv") for i in range(n_files)]

In [7]:
for dict_ in res2:
    print(dict_["30-minutes-or-less"])

{'sum': 278764, 'count': 36602}
{'sum': 276271, 'count': 36499}
{'sum': 279502, 'count': 36712}
{'sum': 282336, 'count': 36790}
{'sum': 280252, 'count': 36692}
{'sum': 275779, 'count': 36405}
{'sum': 275524, 'count': 36307}
{'sum': 279659, 'count': 36724}
{'sum': 278903, 'count': 36613}
{'sum': 276215, 'count': 36438}


3\. Напишите функцию, которая объединяет результаты обработки отдельных файлов. Данная функция принимает на вход список словарей, каждый из которых является результатом вызова функции `get_tag_sum_count_from_file` для конкретного файла, и агрегирует эти словари. Не используйте параллельных вычислений.

Процедура агрегации словарей имеет следующий вид:
$$d_{agg}[k] = \{sum: \sum_{i=1}^{n}d_{i}[k][sum], count: \sum_{i=1}^{n}d_{i}[k][count]\}$$
где $d_1, d_2, ..., d_n$- результат вызова функции `get_tag_sum_count_from_file` для конкретных файлов.

Примените данную функцию к результату выполнения задания 2. Выведите на экран результат для тэга "30-minutes-or-less".

In [8]:
def agg_results(tag_sum_count_list: list) -> dict:
    res = dict()
    main_df = pd.DataFrame(columns=["index", "sum", "count"])
    for dict_ in tag_sum_count_list:
        temp = pd.DataFrame.from_dict(dict_, orient="index").reset_index()
        main_df = pd.concat([temp, main_df])
    return main_df.groupby("index").sum().to_dict("index")

In [9]:
res3 = agg_results(res2)

In [10]:
res3["30-minutes-or-less"]

{'sum': 2783205, 'count': 365782}

4\. Напишите функцию, которая считает среднее значение количества шагов для каждого тэга в словаре, имеющего вид, аналогичный словарям в задаче 2, и возвращает результат в виде словаря . Используйте решения задач 1-3, чтобы получить среднее значение количества шагов каждого тэга для всего датасета, имея результаты обработки частей датасета и результат их агрегации. Выведите на экран результат для тэга "30-minutes-or-less".

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

In [11]:
def get_tag_mean_n_steps(tag_sum_count: dict) -> dict:
    for k,v in tag_sum_count.items():
        tag_sum_count[k] = v["sum"] / v["count"]
    return tag_sum_count

In [12]:
res4 = get_tag_mean_n_steps(agg_results([get_tag_sum_count_from_file(f"tag_nsteps_{i}.csv") for i in range(n_files)]))

In [13]:
res4["30-minutes-or-less"]

7.608917333275011

In [15]:
%%timeit
get_tag_mean_n_steps(agg_results([get_tag_sum_count_from_file(f"tag_nsteps_{i}.csv") for i in range(n_files)]))

14.6 s ± 99 ms per loop (mean ± std. dev. of 7 runs, 1 loop each)


5\. Повторите решение задачи 4, распараллелив вызовы функции `get_tag_sum_count_from_file` для различных файлов с помощью `multiprocessing.Pool`. Для обработки каждого файла создайте свой собственный процесс. Выведите на экран результат для тэга "30-minutes-or-less". Определите, за какое время задача решается для всех файлов. При замере времени учитывайте время расчета статистики для каждого файла, агрегации результатов и, собственно, вычисления средного. Временем, затрачиваемым на процедуру разбиения исходного файла можно пренебречь.

In [14]:
%%file helpers.py

import pandas as pd

def get_tag_sum_count_from_file(file: str) -> dict:
    df = pd.read_csv(file)
    df2 = df.groupby("tags")["n_steps"].agg(sum = 'sum', count = 'count')
    return df2.to_dict('index')

Overwriting helpers.py


In [15]:
import helpers

In [16]:
pool = mp.Pool(processes=n_files)
results = list(pool.map(helpers.get_tag_sum_count_from_file, [f'tag_nsteps_{i}.csv' for i in range(n_files)]))
res5 = get_tag_mean_n_steps(agg_results(results))

In [17]:
res5["30-minutes-or-less"]

7.608917333275011

In [18]:
%%timeit
pool = mp.Pool(processes=n_files)
results = list(pool.map(helpers.get_tag_sum_count_from_file, [f'tag_nsteps_{i}.csv' for i in range(n_files)]))
res5 = get_tag_mean_n_steps(agg_results(results))

3.21 s ± 47.5 ms per loop (mean ± std. dev. of 7 runs, 1 loop each)


6\. Повторите решение задачи 4, распараллелив вычисления функции `get_tag_sum_count_from_file` для различных файлов с помощью `multiprocessing.Process`. Для обработки каждого файла создайте свой собственный процесс. Для обмена данными между процессами используйте `multiprocessing.Queue`.

Выведите на экран результат для тэга "30-minutes-or-less". Определите, за какое время задача решается для всех файлов. При замере времени учитывайте время расчета статистики для каждого файла, агрегации результатов и, собственно, вычисления средного. Временем, затрачиваемым на процедуру разбиения исходного файла можно пренебречь.

In [20]:
%%file helpers2.py
import pandas as pd

def get_tag_sum_count_from_file(file: str, output) -> dict:
    df = pd.read_csv(file)
    df2 = df.groupby("tags")["n_steps"].agg(sum = 'sum', count = 'count')
    output.put(df2.to_dict('index'))

Overwriting helpers2.py


In [21]:
import helpers2

In [22]:
output = mp.Queue()
processes = [mp.Process(target=helpers2.get_tag_sum_count_from_file, args=(f"tag_nsteps_{i}.csv", output)) \
             for i in range(n_files)]
for p in processes:
    p.start()

results = [output.get() for p in processes]
res6 = get_tag_mean_n_steps(agg_results(results))

In [23]:
res6["30-minutes-or-less"]

7.608917333275011

In [24]:
%%timeit
output = mp.Queue()
processes = [mp.Process(target=helpers2.get_tag_sum_count_from_file, args=(f"tag_nsteps_{i}.csv", output)) \
             for i in range(n_files)]
for p in processes:
    p.start()

results = [output.get() for p in processes]
res6 = get_tag_mean_n_steps(agg_results(results))

3.26 s ± 20.7 ms per loop (mean ± std. dev. of 7 runs, 1 loop each)


7\. Исследуйте, как влияет количество запущенных одновременно процессов на скорость решения задачи. Узнайте количество ядер вашего процессора $K$. Повторите решение задачи 1, разбив исходный файл на $\frac{K}{2}$, $K$ и $2K$ фрагментов. Для каждого из разбиений повторите решение задачи 5. Визуализируйте зависимость времени выполнения кода от количества файлов в разбиении. Сделайте вывод в виде текстового комментария.

In [25]:
n_files = 3
files = []
for i in range(n_files):
    files.append(open(f"tag_nsteps_{i}.csv", "w"))
    files[i].write("tags,n_steps\n")
with open("tag_nsteps.csv", "r") as f_main:
    line = f_main.readline()
    line = f_main.readline()
    i = 0
    while line:
        files[i].write(line)
        i = (i + 1) % n_files
        line = f_main.readline()
for i in range(n_files):
    files[i].close()

In [26]:
%%timeit
pool = mp.Pool(processes=n_files)
results = list(pool.map(helpers.get_tag_sum_count_from_file, [f'tag_nsteps_{i}.csv' for i in range(n_files)]))
res5 = get_tag_mean_n_steps(agg_results(results))

5.52 s ± 15.4 ms per loop (mean ± std. dev. of 7 runs, 1 loop each)


In [27]:
n_files = 6
files = []
for i in range(n_files):
    files.append(open(f"tag_nsteps_{i}.csv", "w"))
    files[i].write("tags,n_steps\n")
with open("tag_nsteps.csv", "r") as f_main:
    line = f_main.readline()
    line = f_main.readline()
    i = 0
    while line:
        files[i].write(line)
        i = (i + 1) % n_files
        line = f_main.readline()
for i in range(n_files):
    files[i].close()

In [28]:
%%timeit
pool = mp.Pool(processes=n_files)
results = list(pool.map(helpers.get_tag_sum_count_from_file, [f'tag_nsteps_{i}.csv' for i in range(n_files)]))
res5 = get_tag_mean_n_steps(agg_results(results))

3.56 s ± 86.5 ms per loop (mean ± std. dev. of 7 runs, 1 loop each)


In [29]:
n_files = 12
files = []
for i in range(n_files):
    files.append(open(f"tag_nsteps_{i}.csv", "w"))
    files[i].write("tags,n_steps\n")
with open("tag_nsteps.csv", "r") as f_main:
    line = f_main.readline()
    line = f_main.readline()
    i = 0
    while line:
        files[i].write(line)
        i = (i + 1) % n_files
        line = f_main.readline()
for i in range(n_files):
    files[i].close()

In [30]:
%%timeit
pool = mp.Pool(processes=n_files)
results = list(pool.map(helpers.get_tag_sum_count_from_file, [f'tag_nsteps_{i}.csv' for i in range(n_files)]))
res5 = get_tag_mean_n_steps(agg_results(results))

3.29 s ± 21.9 ms per loop (mean ± std. dev. of 7 runs, 1 loop each)


2K выгодный

8\. Напишите функцию `parallel_map`, которая принимает на вход серию `s` `pd.Series` и функцию одного аргумента `f` и поэлементно применяет эту функцию к серии, распараллелив вычисления при помощи пакета `multiprocessing`. Логика работы функции `parallel_map` должна включать следующие действия:
* разбиение исходной серии на $K$ частей, где $K$ - количество ядер вашего процессора;
* параллельное применение функции `f` к каждой части при помощи метода _серии_ `map` при помощи нескольких подпроцессов;
* объединение результатов работы подпроцессов в одну серию. 

In [1]:
%%file helpers3.py
import pandas as pd

def mapp(s: pd.Series, f: callable) -> pd.Series:
    return s.map(f)

Overwriting helpers3.py


In [2]:
import helpers3

In [3]:
import pandas as pd
import numpy as np
import multiprocessing as mp

In [14]:
s = pd.Series(range(6*100))
def parallel_map(s: pd.Series, f: callable) -> pd.Series: 
    lst = np.array_split(s, 6)
    
    pool = mp.Pool(processes=6)
    return pd.concat(pool.starmap(helpers3.mapp, [(part, f) for part in lst] ))

9\. Напишите функцию `f`, которая принимает на вход тэг и проверяет, удовлетворяет ли тэг следующему шаблону: `[любое число]-[любое слово]-or-less`. Возьмите любой фрагмент файла, полученный в задании 1, примените функцию `f` при помощи `parallel_map` к столбцу `tags` и посчитайте количество тэгов, подходящих под этот шаблон. Решите ту же задачу, воспользовавшись методом _серий_ `map`. Сравните время и результат выполнения двух решений.

In [15]:
%%file helpers4.py
import re
def f(tag: str) -> bool:
    return bool(re.match("\d+-\w+-or-less", str(tag)))

Overwriting helpers4.py


In [16]:
import helpers4

In [17]:
df = pd.read_csv("tag_nsteps_0.csv")

In [19]:
%%time
parallel_map(df.tags, helpers4.f).sum()

CPU times: total: 219 ms
Wall time: 1.46 s


170481

In [20]:
%%time
len(df.tags[df.tags.map(helpers4.f)])

CPU times: total: 3.77 s
Wall time: 3.76 s


170481

10\. Используя пакет `pandarallel`, примените функцию `f` из задания 9 к столбцу `tags` таблицы, с которой вы работали этом задании. Посчитайте количество тэгов, подходящих под описанный шаблон. Измерьте время выполнения кода. Выведите на экран полученный результат.

In [26]:
from pandarallel import pandarallel

In [27]:
pandarallel.initialize()

INFO: Pandarallel will run on 6 workers.
INFO: Pandarallel will use standard multiprocessing data transfer (pipe) to transfer data between the main process and workers.

https://nalepae.github.io/pandarallel/troubleshooting/


In [29]:
%%time
df.tags.parallel_map(helpers4.f).sum()

CPU times: total: 234 ms
Wall time: 1.81 s


170481