
مشکل: 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(). - برای حفظ کارایی، فقط بخشِ دسترسی به منبع مشترک را داخل قفل نگه دارید، نه کل کار سنگین را.