مقدمه
کلاس «Barrier» را در ماژول «threading» دیدیم؛ ابزاری برای متوقف کردن چند ترد تا زمانی که همهی آنها به یک نقطهی مشخص برسند. ماژول «multiprocessing» همین رفتار را برای پردازشها فراهم میکند و کاربرد آن دقیقاً مشابه است.
ایجاد یک Barrier
barrier = multiprocessing.Barrier(3)
عدد «3» تعداد پردازشهایی است که باید به «barrier.wait()» برسند تا همهی آنها همزمان آزاد شوند.
مثال: هماهنگ کردن چند مرحله بین پردازشها
فرض کنید چند پردازش باید کارشان را در چند مرحله انجام دهند و هیچکدام نباید وارد مرحلهی بعد شود مگر اینکه همه به پایان مرحلهی قبل رسیده باشند.
def worker(barrier, worker_id):
print(f'worker-{worker_id}: phase 1 finished')
barrier.wait()
print(f'worker-{worker_id}: phase 2 finished')
barrier.wait()
print(f'worker-{worker_id}: phase 3 finished')
if __name__ == '__main__':
barrier = multiprocessing.Barrier(3)
procs = [
multiprocessing.Process(target=worker, args=(barrier, i))
for i in range(3)
]
for p in procs:
p.start()
for p in procs:
p.join()
در این مثال، هیچ پردازشی وارد «phase 2» نمیشود مگر اینکه هر سه پردازش، «phase 1» را به پایان رسانده و به «barrier.wait()» رسیده باشند. همین قانون برای مرحلهی بعدی هم برقرار است.
وضعیت خراب Barrier
اگر یکی از پردازشها قبل از رسیدن به «barrier.wait()» با خطا متوقف شود یا اصلاً به آن نرسد (مثلاً بهخاطر «timeout»)، «Barrier» وارد حالت «broken» میشود و در همهی پردازشهای دیگری که منتظرند، استثنای «BrokenBarrierError» ایجاد میشود. این رفتار عمداً اینگونه طراحی شده تا پردازشهای باقیمانده برای همیشه منتظر نمانند.
جمعبندی
«Barrier» در ماژول «multiprocessing» همان مفهوم آشنای همزمانسازی چند مرحلهای را به دنیای پردازشها میآورد. استفاده از آن دقیقاً مشابه نسخهی «threading» است، با همان نکتهی همیشگی که شیء باید بین پردازشها به اشتراک گذاشته شود.