Grafy zadań i metody planowania

Programowanie równoległe z Dask w Pythonie

James Fulton

Climate Informatics Researcher

Wizualizacja grafu zadań

# 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()

Diagram pokazuje kroki potrzebne do obliczenia wyniku. my-square-function jest uruchamiana dwukrotnie, a dwa wyjścia są przekazywane do funkcji dodawania, która zwraca jeden wynik.

Programowanie równoległe z Dask w Pythonie

Nakładający się graf zadań

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
Programowanie równoległe z Dask w Pythonie

Nakładający się graf zadań

delayed_result1.visualize()

Diagram przedstawiający graf zadań dla wyniku 1.

delayed_result2.visualize()

Diagram przedstawiający graf zadań dla wyniku 2.

Programowanie równoległe z Dask w Pythonie

Nakładający się graf zadań

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

Graf zadań pokazujący, że wynik 1 i wynik 2 współdzielą wynik pośredni.

Programowanie równoległe z Dask w Pythonie

Wielowątkowość a przetwarzanie równoległe

Przenoszenie danych

Przetwarzanie równoległe
  • Procesy mają własną przestrzeń RAM
Wielowątkowość
  • Wątki współdzielą przestrzeń RAM
Programowanie równoległe z Dask w Pythonie

Wielowątkowość a przetwarzanie równoległe

# 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)
  • Wolne przy przetwarzaniu równoległym

Diagram pokazuje, że dwie tablice z jednego procesu Pythona muszą zostać przesłane do dwóch innych procesów Pythona.

Programowanie równoległe z Dask w Pythonie

Wielowątkowość a przetwarzanie równoległe

# 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)
  • Szybkie przy wielowątkowości

Diagram pokazuje, że dwie tablice nie muszą być w ogóle kopiowane.

Programowanie równoległe z Dask w Pythonie

GIL

Globalny interpreter lock – tylko jeden wątek może odczytywać skrypt Pythona na raz

def sum_to_n(n):
    """Sums numbers from 0 to n"""
    total = 0
    for i in range(n+1):
        total += i
    return total
  • Wielowątkowość tu nie pomoże
  • Przetwarzanie równoległe pomoże
sum1 = delayed(sum_to_n)(1000)
sum2 = delayed(sum_to_n)(1000)
Programowanie równoległe z Dask w Pythonie

Przykładowe czasy – GIL

Trzy wykresy Gantta pokazujące czasy wykonania prostej funkcji Pythona 16 razy. Spośród trzech metod planowania zadań przetwarzanie w procesach było najszybsze.

Programowanie równoległe z Dask w Pythonie

Funkcje zwalniające GIL

  • Np. funkcja pd.read_csv() zwalnia GIL
df1 = delayed(pd.read_csv)('file1.csv')
df2 = delayed(pd.read_csv)('file2.csv')
Programowanie równoległe z Dask w Pythonie

Przykładowe czasy – ładowanie danych

Trzy wykresy Gantta pokazujące czasy ładowania danych z pliku CSV 16 razy. Spośród trzech metod planowania zadań wielowątkowość była najszybsza.

Programowanie równoległe z Dask w Pythonie

Podsumowanie

Wątki

  • Uruchamiają się bardzo szybko
  • Współdzielą przestrzeń pamięci z główną sesją
  • Nie wymagają transferu pamięci
  • Ograniczone przez GIL, który pozwala jednemu wątkowi odczytywać kod na raz

Procesy

  • Wymagają czasu i pamięci do uruchomienia
  • Mają oddzielne pule pamięci
  • Transfer danych między nimi a główną sesją Pythona jest bardzo wolny
  • Każdy ma własny GIL i nie musi czekać na odczyt kodu
Programowanie równoległe z Dask w Pythonie

Czas na ćwiczenia!

Programowanie równoległe z Dask w Pythonie

Preparing Video For Download...