任務圖與排程方法

在 Python 中使用 Dask 進行平行程式設計

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 執行兩次,兩個輸出傳入 add 函式,回傳一個輸出。

在 Python 中使用 Dask 進行平行程式設計

重疊的任務圖

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
在 Python 中使用 Dask 進行平行程式設計

重疊的任務圖

delayed_result1.visualize()

顯示結果 1 任務圖的示意圖。

delayed_result2.visualize()

顯示結果 2 任務圖的示意圖。

在 Python 中使用 Dask 進行平行程式設計

重疊的任務圖

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

一張任務圖,顯示結果 1 與結果 2 共用一個中介結果。

在 Python 中使用 Dask 進行平行程式設計

多執行緒 vs. 平行處理

資料移動

平行處理
  • 行程各自有獨立 RAM 空間
多執行緒
  • 執行緒共用同一個 RAM 空間
在 Python 中使用 Dask 進行平行程式設計

多執行緒 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 行程。

在 Python 中使用 Dask 進行平行程式設計

多執行緒 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)
  • 用多執行緒很快

示意圖:兩個陣列完全不需要被複製。

在 Python 中使用 Dask 進行平行程式設計

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)
在 Python 中使用 Dask 進行平行程式設計

範例時間比較-GIL

三個甘特圖,顯示將簡單 Python 函式執行 16 次的時間。三種排程方法中,行程(processes)最快。

在 Python 中使用 Dask 進行平行程式設計

會釋放 GIL 的函式

  • 例如,pd.read_csv() 會釋放 GIL
df1 = delayed(pd.read_csv)('file1.csv')
df2 = delayed(pd.read_csv)('file2.csv')
在 Python 中使用 Dask 進行平行程式設計

範例時間比較-載入資料

三個甘特圖,顯示將載入 CSV 的函式執行 16 次的時間。三種排程方法中,執行緒(threads)最快。

在 Python 中使用 Dask 進行平行程式設計

重點整理

執行緒

  • 啟動非常快
  • 與主工作階段共用記憶體空間
  • 不需傳輸記憶體
  • 受 GIL 限制,一次只允許一個執行緒讀取程式碼

行程

  • 建立耗費時間與記憶體
  • 各自有獨立記憶體池
  • 與彼此及主 Python 工作階段之間的資料傳輸很慢
  • 各自有 GIL,無需輪流讀取程式碼
在 Python 中使用 Dask 進行平行程式設計

一起來練習吧!

在 Python 中使用 Dask 進行平行程式設計

Preparing Video For Download...