กราฟงานและวิธีการจัดตารางงาน

Parallel Programming with Dask in 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 ทำงาน 2 ครั้ง แล้วส่งผลลัพธ์ทั้งสองไปยังฟังก์ชัน add ซึ่งคืนค่าผลลัพธ์หนึ่งค่า

Parallel Programming with Dask in 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
Parallel Programming with Dask in Python

กราฟงานที่ซ้อนทับกัน

delayed_result1.visualize()

แผนภาพแสดงกราฟงานของ result 1

delayed_result2.visualize()

แผนภาพแสดงกราฟงานของ result 2

Parallel Programming with Dask in Python

กราฟงานที่ซ้อนทับกัน

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

กราฟงานที่แสดงว่า result 1 และ result 2 ใช้ผลลัพธ์ขั้นกลางร่วมกัน

Parallel Programming with Dask in Python

Multi-threading กับการประมวลผลแบบขนาน

การย้ายข้อมูล

การประมวลผลแบบขนาน
  • แต่ละ process มีพื้นที่ RAM เป็นของตัวเอง
Multi-threading
  • Thread ใช้พื้นที่ RAM ร่วมกัน
Parallel Programming with Dask in Python

Multi-threading กับการประมวลผลแบบขนาน

# 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 process เดียวกันต้องถูกส่งไปยัง process อื่นอีก 2 ตัว

Parallel Programming with Dask in Python

Multi-threading กับการประมวลผลแบบขนาน

# 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)
  • เร็วเมื่อใช้ multi-threading

แผนภาพแสดงว่าไม่จำเป็นต้องคัดลอกอาร์เรย์ทั้งสองเลย

Parallel Programming with Dask in Python

GIL

Global interpreter lock — อนุญาตให้ thread เดียวเท่านั้นอ่าน Python script ได้ในแต่ละขณะ

def sum_to_n(n):
    """Sums numbers from 0 to n"""
    total = 0
    for i in range(n+1):
        total += i
    return total
  • Multi-threading ไม่ช่วยในกรณีนี้
  • การประมวลผลแบบขนานช่วยได้
sum1 = delayed(sum_to_n)(1000)
sum2 = delayed(sum_to_n)(1000)
Parallel Programming with Dask in Python

ตัวอย่างเวลาการรัน — GIL

แผนภูมิ Gantt 3 แผ่นที่แสดงเวลาการรันฟังก์ชัน Python แบบง่าย 16 ครั้ง จากวิธีการจัดตารางงาน 3 แบบ พบว่า processes รันได้เร็วที่สุด

Parallel Programming with Dask in Python

ฟังก์ชันที่ปล่อย GIL

  • เช่น ฟังก์ชัน pd.read_csv() จะปล่อย GIL
df1 = delayed(pd.read_csv)('file1.csv')
df2 = delayed(pd.read_csv)('file2.csv')
Parallel Programming with Dask in Python

ตัวอย่างเวลาการรัน — การโหลดข้อมูล

แผนภูมิ Gantt 3 แผ่นที่แสดงเวลาการรันฟังก์ชันโหลดข้อมูลจาก CSV 16 ครั้ง จากวิธีการจัดตารางงาน 3 แบบ พบว่า threads รันได้เร็วที่สุด

Parallel Programming with Dask in Python

สรุป

Threads

  • เริ่มต้นทำงานได้รวดเร็วมาก
  • ใช้พื้นที่หน่วยความจำร่วมกับ session หลัก
  • ไม่ต้องถ่ายโอนหน่วยความจำ
  • ถูกจำกัดโดย GIL ซึ่งอนุญาตให้ thread เดียวอ่านโค้ดได้ในแต่ละขณะ

Processes

  • ใช้เวลาและหน่วยความจำในการตั้งค่า
  • มีพื้นที่หน่วยความจำแยกจากกัน
  • การถ่ายโอนข้อมูลระหว่าง process และไปยัง Python session หลักทำได้ช้ามาก
  • แต่ละ process มี GIL ของตัวเอง จึงไม่ต้องสลับกันอ่านโค้ด
Parallel Programming with Dask in Python

มาฝึกกันเถอะ!

Parallel Programming with Dask in Python

Preparing Video For Download...