آموزش ماژول multiprocessing – کلاس Lock

Please login to bookmark Close

مشکل: Race Condition بین پردازه‌ها

وقتی چند پردازه به یک منبع مشترک (مثل یک Value) دسترسی هم‌زمان دارند و آن را تغییر می‌دهند، همان مشکل Race Condition که در تردها دیدیم اینجا هم رخ می‌دهد. به کد زیر توجه کنید که پنج پردازه هرکدام ۱۰۰٬۰۰۰ بار یک شمارنده‌ی مشترک را زیاد می‌کنند:

import multiprocessing

# shared resource
shared_counter = multiprocessing.Value('i', 0)

def increment_counter(shared_counter):
    for _ in range(100000):
        shared_counter.value += 1

if __name__ == '__main__':
    processes = []
    for _ in range(5):
        p = multiprocessing.Process(target=increment_counter, args=(shared_counter,))
        p.start()
        processes.append(p)

    for p in processes:
        p.join()

    print("Final value of the counter:", shared_counter.value)

انتظار داریم خروجی 500000 باشد، اما مقداری کمتر (و هر بار متفاوت) می‌گیریم. دلیلش این است که shared_counter.value += 1 یک عمل اتمیک نیست؛ سه گامِ خواندن، افزودن و نوشتن دارد و چند پردازه می‌توانند وسط این گام‌ها یکدیگر را بازنویسی کنند.

راه‌حل: multiprocessing.Lock

برای حل مشکل از کلاس Lock ماژول multiprocessing استفاده می‌کنیم تا در هر لحظه فقط یک پردازه وارد بخش بحرانی شود. قفل را می‌سازیم و به هر پردازه پاس می‌دهیم؛ پیش از کار با منبع مشترک acquire و پس از آن release می‌کنیم:

import multiprocessing

lock = multiprocessing.Lock()
shared_counter = multiprocessing.Value('i', 0)

def increment_counter(shared_counter, lock):
    lock.acquire()
    for _ in range(100000):
        shared_counter.value += 1
    lock.release()

if __name__ == '__main__':
    processes = []
    for _ in range(5):
        p = multiprocessing.Process(target=increment_counter, args=(shared_counter, lock))
        p.start()
        processes.append(p)

    for p in processes:
        p.join()

    print("Final value of the counter:", shared_counter.value)   # 500000

حالا خروجی همیشه 500000 است. توجه کنید که قفل باید در پردازه‌ی اصلی ساخته و به پردازه‌ها پاس داده شود (نه اینکه هر پردازه قفل خودش را بسازد)، وگرنه اشتراکی نخواهد بود.

روش تمیزتر: with

مثل نسخه‌ی تردی، اینجا هم می‌توان به‌جای acquire/release دستی از with استفاده کرد تا قفل حتی در صورت بروز خطا هم آزاد شود:

import multiprocessing

lock = multiprocessing.Lock()
shared_counter = multiprocessing.Value('i', 0)

def increment_counter(shared_counter, lock):
    with lock:
        for _ in range(100000):
            shared_counter.value += 1

if __name__ == '__main__':
    processes = []
    for _ in range(5):
        p = multiprocessing.Process(target=increment_counter, args=(shared_counter, lock))
        p.start()
        processes.append(p)

    for p in processes:
        p.join()

    print("Final value of the counter:", shared_counter.value)

نکته: قفل داخلی خودِ Value

در واقع خودِ Value یک قفل داخلی دارد که با shared_counter.get_lock() در دسترس است؛ پس به‌جای ساختن یک قفل جداگانه می‌توانید از همان استفاده کنید:

def increment_counter(shared_counter):
    for _ in range(100000):
        with shared_counter.get_lock():   # Value has its own built-in lock
            shared_counter.value += 1

توجه به کارایی: قفل‌کردنِ کل حلقه (مثل مثال‌های بالا) درست است اما پردازه‌ها را عملاً سریالی می‌کند و از موازی‌سازی چیزی باقی نمی‌ماند. اگر کار واقعیِ سنگینی خارج از بخش بحرانی دارید، فقط همان بخشی را که به منبع مشترک دست می‌زند داخل قفل بگذارید، نه کل کار را.

مثال واقعی: شمارش امنِ نتایج بین کارگرها

فرض کنید چند پردازه‌ی کارگر دارید که هرکدام بخشی از داده را پردازش می‌کنند و می‌خواهند تعداد کل موارد موفق را در یک شمارنده‌ی مشترک جمع بزنند. با قفل مطمئن می‌شویم این جمع‌زدن درست انجام شود:

import multiprocessing

def worker(chunk, success_count, lock):
    local = sum(1 for x in chunk if x % 2 == 0)   # count even numbers (the "work")
    with lock:
        success_count.value += local               # safely add to the shared total

if __name__ == '__main__':
    lock = multiprocessing.Lock()
    success_count = multiprocessing.Value('i', 0)

    data = list(range(100))
    chunks = [data[i::4] for i in range(4)]         # split data into 4 parts

    procs = [multiprocessing.Process(target=worker, args=(c, success_count, lock)) for c in chunks]
    for p in procs: p.start()
    for p in procs: p.join()

    print("total evens found:", success_count.value)   # 50

هر کارگر ابتدا کار سنگینش را بیرون از قفل انجام می‌دهد و فقط برای «جمع‌زدنِ نتیجه در شمارنده‌ی مشترک» قفل می‌گیرد؛ این‌طور هم درستی حفظ می‌شود و هم موازی‌سازی تا حد ممکن باقی می‌ماند.

جمع‌بندی

  • دسترسی هم‌زمان چند پردازه به یک Value/Array مشترک می‌تواند Race Condition ایجاد کند، چون += اتمیک نیست.
  • با multiprocessing.Lock بخش بحرانی را محافظت می‌کنیم؛ قفل را در پردازه‌ی اصلی بسازید و به پردازه‌ها پاس دهید.
  • به‌جای acquire/release از with استفاده کنید؛ یا از قفل داخلیِ Value.get_lock().
  • برای حفظ کارایی، فقط بخشِ دسترسی به منبع مشترک را داخل قفل نگه دارید، نه کل کار سنگین را.

Please login to bookmark Close
پیشرفت شما در «دوره آموزش کانکارنسی در پایتون» (76%)
نظرات

دیدگاهتان را بنویسید

76%
پیشرفت

سرفصل دوره

فهرست مطالب

سرفصل دوره

تمرین

این قسمت تمرین ندارد!

پاسخ تمرین ها

هنوز برای تمرین‌های این قسمت پاسخی ثبت نشده است!

اشتراک گذاری

چرا بهتره از فیلترشکن استفاده کنید؟

من همه ویدئو ها و پادکست های کُدباز رو توی یوتیوب و ساندکلود و پلتفرم هایی آپلود می‌کنم که اغلب فیلتر هستند.

اغلب آموزش‌ها ویدئو و پادکست دارند. پس اگر می‌خواهید از محتوای سایت بیشترین استفاده رو ببرید نیاز به فیلتر شکن دارید.

توجه داشته باشید که برای خرید از فروشگاه بهتره فیلتر شکن رو خاموش کنید.

تنظیمات

انتخاب زبان
تغییر تم