Графи завдань і методи планування

Паралельне програмування з Dask у Python

James Fulton

Climate Informatics Researcher

Візуалізація графа завдань

# Create 2 delayed objects
delayed_num1 = delayed(my_square_function)(3)
delayed_num2 = delayed(my_square_function)(4)

# Add them
result = delayed_num1 + delayed_num2

# Plot the task graph result.visualize()

Діаграма кроків для обчислення результату: my-square-function виконується двічі, обидва виводи передаються у функцію додавання, що повертає один результат.

Паралельне програмування з Dask у Python

Перекривний граф завдань

delayed_intermediate = delayed(my_square_function)(3)

# These two results both use delayed_intermediate_result
delayed_result1 = delayed_intermediate - 5
delayed_result2 = delayed_intermediate + 4
Паралельне програмування з Dask у Python

Перекривний граф завдань

delayed_result1.visualize()

Діаграма графа завдань для результату 1.

delayed_result2.visualize()

Діаграма графа завдань для результату 2.

Паралельне програмування з Dask у Python

Перекривний граф завдань

# Plot the task graph
dask.visualize(delayed_result1, delayed_result2)

Граф завдань, який показує, що результат 1 і результат 2 мають спільний проміжний результат.

Паралельне програмування з Dask у Python

Мульти-трединг vs паралельна обробка

Переміщення даних

Паралельна обробка
  • Процеси мають власний простір RAM
Мульти-трединг
  • Потоки використовують спільний простір RAM
Паралельне програмування з Dask у Python

Мульти-трединг vs паралельна обробка

# Run a sum on two big arrays
sum1 = delayed(np.sum)(big_array1)
sum2 = delayed(np.sum)(big_array2)

# Compute using processes
dask.compute(sum1, sum2)
  • Повільно з паралельною обробкою

Діаграма показує, що два масиви, створені в одному процесі Python, потрібно надсилати до двох інших процесів Python.

Паралельне програмування з Dask у Python

Мульти-трединг vs паралельна обробка

# Run a sum on two big arrays
sum1 = delayed(np.sum)(big_array1)
sum2 = delayed(np.sum)(big_array2)

# Compute using threads
dask.compute(sum1, sum2)
  • Швидко з мульти-тредингом

Діаграма показує, що два масиви взагалі не потрібно копіювати.

Паралельне програмування з Dask у Python

GIL

Global interpreter lock — лише один потік може читати скрипт Python одночасно

def sum_to_n(n):
    """Sums numbers from 0 to n"""
    total = 0
    for i in range(n+1):
        total += i
    return total
  • Мульти-трединг тут не допоможе
  • Паралельна обробка — так
sum1 = delayed(sum_to_n)(1000)
sum2 = delayed(sum_to_n)(1000)
Паралельне програмування з Dask у Python

Приклад часу — GIL

Три діаграми Ґанта з часами виконання простої Python-функції 16 разів. Із трьох методів планування найшвидшими виявилися процеси.

Паралельне програмування з Dask у Python

Функції, що звільняють GIL

  • Напр., функція pd.read_csv() звільняє GIL
df1 = delayed(pd.read_csv)('file1.csv')
df2 = delayed(pd.read_csv)('file2.csv')
Паралельне програмування з Dask у Python

Приклад часу — завантаження даних

Три діаграми Ґанта з часами виконання функції завантаження CSV 16 разів. Із трьох методів планування найшвидшими були потоки.

Паралельне програмування з Dask у Python

Підсумок

Потоки

  • Дуже швидко запускаються
  • Спільний простір пам'яті з основною сесією
  • Передача пам'яті не потрібна
  • Обмежені GIL: один потік читає код одночасно

Процеси

  • Потребують часу й пам'яті для налаштування
  • Окремі пули пам'яті
  • Дуже повільна передача даних між собою та в основну сесію Python
  • Кожен має власний GIL, тож не треба чергувати читання коду
Паралельне програмування з Dask у Python

Давайте потренуємось!

Паралельне програмування з Dask у Python

Preparing Video For Download...