multiprocessing --- موازیسازی مبتنی بر فرایند¶
کد منبع: Lib/multiprocessing/
دسترسپذیری: not Android, not iOS, not WASI.
این ماژول در پلتفرمهای موبایل یا پلتفرمهای WebAssembly پشتیبانی نمیشود.
مقدمه¶
multiprocessing بستهای است که از ایجاد فرایندها با استفاده از یک API مشابه ماژول threading پشتیبانی میکند. بسته multiprocessing همروندی محلی و راهدور را ارائه میدهد و عملاً با استفاده از زیرفرایندها بهجای نخها، قفل مفسر سراسری را دور میزند. به همین دلیل، ماژول multiprocessing به برنامهنویس اجازه میدهد تا بهطور کامل از چندین پردازنده روی یک ماشین مشخص بهره ببرد. این ماژول روی هر دو POSIX و Windows اجرا میشود.
ماژول multiprocessing همچنین شیء Pool را معرفی میکند که وسیلهای مناسب برای موازیسازی اجرای یک تابع روی چندین مقدار ورودی فراهم میآورد و دادههای ورودی را بین فرآیندها توزیع میکند (موازیسازی دادهای). مثال زیر رویه رایج تعریف چنین توابعی در یک ماژول را نشان میدهد تا فرآیندهای فرزند بتوانند آن ماژول را با موفقیت ایمپورت کنند. این مثال پایه از موازیسازی دادهای با استفاده از Pool،
from multiprocessing import Pool
def f(x):
return x*x
if __name__ == '__main__':
with Pool(5) as p:
print(p.map(f, [1, 2, 3]))
به خروجی استاندارد چاپ میشود
[1, 4, 9]
ماژول multiprocessing همچنین APIهایی را معرفی میکند که معادلی در ماژول threading ندارند، مانند توانایی terminate، interrupt یا kill یک فرآیند در حال اجرا.
همچنین ملاحظه نمائید
concurrent.futures.ProcessPoolExecutor یک رابط سطح بالاتر برای ارسال وظایف به یک فرایند پسزمینه بدون مسدود کردن اجرای فرایند فراخوان ارائه میدهد. در مقایسه با استفاده مستقیم از رابط Pool، API concurrent.futures بهراحتی بیشتری امکان میدهد که ارسال کار به استخر فرایند زیربنایی از انتظار برای نتایج جدا شود.
کلاس Process¶
در multiprocessing، فرایندها با ایجاد یک شیء Process و سپس فراخوانی متد start() آن به وجود میآیند. Process از APIِ threading.Thread پیروی میکند. یک مثال ساده از یک برنامه چندفرایندی به صورت زیر است
from multiprocessing import Process
def f(name):
print('hello', name)
if __name__ == '__main__':
p = Process(target=f, args=('bob',))
p.start()
p.join()
برای نمایش شناسههای فرایند درگیر بهتفکیک، در اینجا یک مثال گسترشیافته آمده است:
from multiprocessing import Process
import os
def info(title):
print(title)
print('module name:', __name__)
print('parent process:', os.getppid())
print('process id:', os.getpid())
def f(name):
info('function f')
print('hello', name)
if __name__ == '__main__':
info('main line')
p = Process(target=f, args=('bob',))
p.start()
p.join()
برای توضیحی در مورد اینکه چرا بخش if __name__ == '__main__' ضروری است، دستورالعملهای برنامهنویسی را ببینید.
آرگومانهای Process معمولاً باید پیکلپذیر باشند تا بتوان آنها را به فرایند فرزند منتقل کرد. اگر مثال بالا را مستقیماً در یک REPL تایپ کنید، ممکن است در فرایند فرزند، هنگام تلاش برای یافتن تابع f در ماژول __main__، یک AttributeError رخ دهد.
زمینهها و متدهای شروع¶
بسته به سکو، multiprocessing از سه روش برای شروع یک فرایند پشتیبانی میکند. این روشهای شروع عبارتند از
- spawn
فرایند والد، یک فرایند مفسر پایتون تازه را آغاز میکند. فرایند فرزند تنها منابع مورد نیاز برای اجرای متد
run()شیء فرایند را به ارث میبرد. بهطور خاص، توصیفگرهای پرونده و دستههای غیرضروری از فرایند والد به ارث برده نمیشوند. آغاز یک فرایند با استفاده از این روش، در مقایسه با استفاده از fork یا forkserver نسبتاً کند است.در دسترس در سکوهای POSIX و Windows. پیشفرض در Windows و macOS.
- fork
فرایند والد از
os.fork()برای فورک کردن مفسر پایتون استفاده میکند. فرایند فرزند، هنگامی که آغاز میشود، عملاً با فرایند والد یکسان است. فرایند فرزند تمام منابع فرایند والد را به ارث میبرد. توجه داشته باشید که فورک کردن بهصورت ایمن یک فرایند چندنخی مشکلساز است.در سیستمهای POSIX در دسترس است.
تغییر یافته در نسخهی 3.14: این دیگر روش پیشفرض شروع در هیچ سکویی نیست. کدی که به fork نیاز دارد، باید صریحاً آن را از طریق
get_context()یاset_start_method()مشخص کند.تغییر یافته در نسخهی 3.12: اگر پایتون بتواند تشخیص دهد که فرایند شما دارای چندین نخ است، تابع
os.fork()که این متد شروع آن را بهصورت داخلی فراخوانی میکند، یکDeprecationWarningرا پرتاب میکند. از یک متد شروع دیگر استفاده کنید. برای توضیح بیشتر، مستنداتos.fork()را ببینید.
- forkserver
هنگامی که برنامه شروع میشود و متد شروع forkserver را انتخاب میکند، یک فرایند سرور ایجاد میشود. از آن پس، هر زمان که به یک فرایند جدید نیاز باشد، فرایند والد به سرور متصل میشود و از آن درخواست میکند که یک فرایند جدید را fork کند. فرایند سرور fork تکنخی است، مگر اینکه کتابخانههای سیستمی یا ایمپورتهای از پیش بارگذاریشده بهعنوان یک اثر جانبی نخهایی را ایجاد کنند؛ بنابراین استفاده از
os.fork()برای آن عموماً امن است. هیچ منبع غیرضروری به ارث برده نمیشود.در پلتفرمهای POSIX مانند لینوکس که از انتقال توصیفگرهای پرونده از طریق پایپ یونیکس پشتیبانی میکنند، در دسترس است. پیشفرض در آنها.
تغییر یافته در نسخهی 3.14: این به روش شروع پیشفرض در سکوهای POSIX تبدیل شد.
تغییر یافته در نسخهی 3.4: spawn در همهی پلتفرمهای POSIX اضافه شد، و forkserver برای برخی پلتفرمهای POSIX افزوده شد. فرآیندهای فرزند دیگر همهی دستههای قابل ارثبری والد را در ویندوز به ارث نمیبرند.
تغییر یافته در نسخهی 3.8: در macOS، روش شروع spawn اکنون پیشفرض است. روش شروع fork باید ناامن در نظر گرفته شود، زیرا میتواند باعث خرابی زیرفرایند شود، چون کتابخانههای سیستم macOS ممکن است نخها را شروع کنند. به bpo-33725 مراجعه کنید.
تغییر یافته در نسخهی 3.14: در پلتفرمهای POSIX، روش شروع پیشفرض از fork به forkserver تغییر کرد تا ضمن حفظ کارایی، از ناسازگاریهای رایج فرآیندهای چندنخی اجتناب شود. gh-84559 را ببینید.
در POSIX، استفاده از روشهای شروع spawn یا forkserver یک فرایند ردیاب منابع (resource tracker) را نیز راهاندازی میکند که منابع سیستمی نامداری را که توسط فرایندهای برنامه ایجاد شدهاند و پیوندشان حذف شده است (unlinked) (مانند سمافورهای نامدار یا اشیای SharedMemory)، پیگیری میکند. پس از خروج همهی فرایندها، ردیاب منابع پیوند هر شیء پیگیریشدهی باقیمانده را حذف میکند. معمولاً نباید هیچکدام وجود داشته باشد، اما اگر فرایندی با یک سیگنال کشته شود، ممکن است برخی منابع «نشتکرده» وجود داشته باشند. (نه سمافورهای نشتکرده و نه بخشهای حافظهی مشترک، بهطور خودکار تا راهاندازی مجدد بعدی پیوندشان حذف نخواهد شد. این موضوع برای هر دو شیء مشکلساز است، زیرا سیستم فقط تعداد محدودی سمافور نامدار را مجاز میداند و بخشهای حافظهی مشترک مقداری فضا در حافظهی اصلی اشغال میکنند.)
برای انتخاب یک روش شروع، از set_start_method() در بند if __name__ == '__main__' ماژول اصلی استفاده میکنید. برای مثال:
import multiprocessing as mp
def foo(q):
q.put('hello')
if __name__ == '__main__':
mp.set_start_method('spawn')
q = mp.Queue()
p = mp.Process(target=foo, args=(q,))
p.start()
print(q.get())
p.join()
set_start_method() نباید بیش از یکبار در برنامه استفاده شود.
بهعنوان جایگزین، میتوانید از get_context() برای به دست آوردن یک شیء زمینه استفاده کنید. اشیاء زمینه همان API ماژول multiprocessing را دارند و امکان استفاده از چندین روش شروع در یک برنامه را فراهم میکنند.
import multiprocessing as mp
def foo(q):
q.put('hello')
if __name__ == '__main__':
ctx = mp.get_context('spawn')
q = ctx.Queue()
p = ctx.Process(target=foo, args=(q,))
p.start()
print(q.get())
p.join()
توجه داشته باشید که اشیای مرتبط با یک زمینه ممکن است با فرآیندهای یک زمینهی دیگر سازگار نباشند. بهویژه، قفلهای ایجادشده با استفاده از زمینهی fork نمیتوانند به فرآیندهایی که با استفاده از روشهای شروع spawn یا forkserver آغاز شدهاند، منتقل شوند.
کتابخانههایی که از multiprocessing یا ProcessPoolExecutor استفاده میکنند، باید بهگونهای طراحی شوند که به کاربران خود امکان دهند زمینهی چندپردازشی خود را فراهم کنند. استفاده از یک زمینهی خاص خودتان در یک کتابخانه میتواند باعث ناسازگاری با سایر بخشهای برنامهی کاربر کتابخانه شود. همیشه در صورتی که کتابخانه شما به روش شروع خاصی نیاز دارد، آن را مستند کنید.
هشدار
بهطور کلی، نمیتوان از روشهای شروع 'spawn' و 'forkserver' با پروندههای اجرایی «فریزشده» (یعنی پروندههای دودویی تولیدشده توسط بستههایی مانند PyInstaller و cx_Freeze) در سیستمهای POSIX استفاده کرد. روش شروع 'fork' ممکن است در صورتی کار کند که کد از نخها استفاده نکند.
مبادلهی شیءها بین فرایندها¶
multiprocessing از دو نوع کانال ارتباطی بین فرآیندها پشتیبانی میکند:
صفها
کلاس
Queueتقریباً یک کپی ازqueue.Queueاست. برای مثال:from multiprocessing import Process, Queue def f(q): q.put([42, None, 'hello']) if __name__ == '__main__': q = Queue() p = Process(target=f, args=(q,)) p.start() print(q.get()) # prints "[42, None, 'hello']" p.join()صفها برای نخ و فرایند ایمن هستند. هر شیء قرار دادهشده در یک صف
multiprocessingسریالسازی خواهد شد.
پایپها
تابع
Pipe()یک جفت شیء اتصال را برمیگرداند که از طریق یک پایپ به هم متصل شدهاند و این پایپ بهطور پیشفرض دوطرفه (دوسویه) است. برای مثال:from multiprocessing import Process, Pipe def f(conn): conn.send([42, None, 'hello']) conn.close() if __name__ == '__main__': parent_conn, child_conn = Pipe() p = Process(target=f, args=(child_conn,)) p.start() print(parent_conn.recv()) # prints "[42, None, 'hello']" p.join()دو شیء اتصال برگرداندهشده توسط
Pipe()، نشاندهنده دو سر پایپ هستند. هر شیء اتصال دارای متدهایsend()وrecv()(در میان سایر متدها) است. توجه داشته باشید که اگر دو فرآیند (یا نخ) تلاش کنند همزمان از سر یکسان پایپ بخوانند یا در آن بنویسند، ممکن است دادههای درون پایپ خراب شوند. البته هیچ خطر خرابی ناشی از فرآیندهایی که همزمان از سرهای مختلف پایپ استفاده میکنند، وجود ندارد.متد
send()شیء را سریالسازی میکند وrecv()شیء را بازسازی میکند.
همگامسازی میان فرآیندها¶
multiprocessing شامل معادلهایی برای تمام سازوکارهای اولیهی همگامسازی در threading است. برای مثال، میتوان از یک قفل استفاده کرد تا اطمینان حاصل شود که تنها یک فرایند در هر لحظه روی خروجی استاندارد چاپ میکند:
from multiprocessing import Process, Lock
def f(l, i):
l.acquire()
try:
print('hello world', i)
finally:
l.release()
if __name__ == '__main__':
lock = Lock()
for num in range(10):
Process(target=f, args=(lock, num)).start()
بدون استفاده از قفل، خروجی فرایندهای مختلف ممکن است کاملاً درهم شود.
استفاده از استخری از کارگرها¶
کلاس Pool نمایانگر یک استخر از فرایندهای کارگر است. این کلاس دارای متدهایی است که امکان واگذاری وظایف به فرایندهای کارگر را به چند روش مختلف فراهم میکنند.
برای مثال:
from multiprocessing import Pool, TimeoutError
import time
import os
def f(x):
return x*x
if __name__ == '__main__':
# start 4 worker processes
with Pool(processes=4) as pool:
# print "[0, 1, 4,..., 81]"
print(pool.map(f, range(10)))
# print same numbers in arbitrary order
for i in pool.imap_unordered(f, range(10)):
print(i)
# evaluate "f(20)" asynchronously
res = pool.apply_async(f, (20,)) # runs in *only* one process
print(res.get(timeout=1)) # prints "400"
# evaluate "os.getpid()" asynchronously
res = pool.apply_async(os.getpid, ()) # runs in *only* one process
print(res.get(timeout=1)) # prints the PID of that process
# launching multiple evaluations asynchronously *may* use more processes
multiple_results = [pool.apply_async(os.getpid, ()) for i in range(4)]
print([res.get(timeout=1) for res in multiple_results])
# make a single worker sleep for 10 seconds
res = pool.apply_async(time.sleep, (10,))
try:
print(res.get(timeout=1))
except TimeoutError:
print("We lacked patience and got a multiprocessing.TimeoutError")
print("For the moment, the pool remains available for more work")
# exiting the 'with'-block has stopped the pool
print("Now the pool is closed and no longer available")
توجه داشته باشید که متدهای یک استخر فقط باید توسط فرایندی که آن را ایجاد کرده است استفاده شوند.
توجه
کارکرد درون این بسته مستلزم آن است که فرزندان بتوانند ماژول __main__ را ایمپورت کنند. این موضوع در دستورالعملهای برنامهنویسی پوشش داده شده است، اما ارزش دارد که در اینجا به آن اشاره شود. این بدان معناست که برخی مثالها، مانند مثالهای multiprocessing.pool.Pool در مفسر تعاملی کار نخواهند کرد. برای مثال:
>>> from multiprocessing import Pool
>>> p = Pool(5)
>>> def f(x):
... return x*x
...
>>> with p:
... p.map(f, [1,2,3])
Process PoolWorker-1:
Process PoolWorker-2:
Process PoolWorker-3:
Traceback (most recent call last):
Traceback (most recent call last):
Traceback (most recent call last):
AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>
AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>
AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>
(اگر این را امتحان کنید، در واقع سه ردگیری کامل پشته را بهصورت نیمهتصادفی درهمآمیخته خروجی میدهد، و سپس ممکن است مجبور شوید بهنحوی فرایند والد را متوقف کنید.)
مرجع¶
بستهی multiprocessing عمدتاً API ماژول threading را بازتولید میکند.
روش شروع سراسری¶
پایتون از چندین روش برای ایجاد و مقداردهی اولیهی یک فرایند پشتیبانی میکند. متد آغاز سراسری، سازوکار پیشفرض برای ایجاد یک فرایند را تنظیم میکند.
چندین تابع و متد چندپردازشی (multiprocessing) که ممکن است اشیاء خاصی را نیز نمونهسازی کنند، اگر روش شروع سراسری پیشتر تنظیم نشده باشد، آن را بهطور ضمنی روی پیشفرض سیستم تنظیم میکنند. روش شروع سراسری فقط یک بار میتواند تنظیم شود. اگر نیاز دارید روش شروع را از پیشفرض سیستم تغییر دهید، باید پیش از فراخوانی توابع یا متدها، یا ایجاد این اشیاء، روش شروع سراسری را از پیش تنظیم کنید.
Process و استثناها¶
- class multiprocessing.Process(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)¶
اشیای Process نشاندهنده فعالیتی هستند که در یک فرایند جداگانه اجرا میشود. کلاس
Processمعادلهایی برای همه متدهایthreading.Threadدارد.سازنده باید همیشه با آرگومانهای کلیدواژهای فراخوانی شود. group باید همیشه
Noneباشد؛ این صرفاً برای سازگاری باthreading.Threadوجود دارد. target شیء فراخوانیپذیر است که توسط متدrun()فراخوانی میشود. مقدار پیشفرض آنNoneاست، یعنی هیچ چیزی فراخوانی نمیشود. name نام فرایند است (برای جزئیات بیشترnameرا ببینید). args تاپل آرگومانهای فراخوانی هدف است. kwargs دیکشنری از آرگومانهای کلیدواژهای برای فراخوانی هدف است. در صورت ارائه، آرگومان daemon که فقط کلیدواژهای است، پرچمdaemonفرایند را رویTrueیاFalseتنظیم میکند. اگرNone(پیشفرض) باشد، این پرچم از فرایند ایجادکننده به ارث میرسد.بهطور پیشفرض، هیچ آرگومانی به target ارسال نمیشود. میتوان از آرگومان args، که مقدار پیشفرض آن
()است، برای مشخصکردن فهرستی یا تاپلی از آرگومانها جهت ارسال به target استفاده کرد.اگر یک زیرکلاس سازنده را بازنویسی کند، باید اطمینان حاصل کند که سازنده کلاس پایه (
super().__init__()) را پیش از انجام هر کار دیگری روی فرایند فراخوانی میکند.توجه
بهطور کلی، تمام آرگومانهای
Processباید پیکلپذیر باشند. این موضوع اغلب هنگام تلاش برای ایجاد یکProcessیا استفاده ازconcurrent.futures.ProcessPoolExecutorاز یک پوستهی تعاملی (REPL) با یک تابع target تعریفشده بهصورت محلی مشاهده میشود.ارسال یک شیء فراخوانیپذیر تعریفشده در نشست REPL فعلی، باعث میشود فرایند فرزند هنگام شروع در اثر یک استثنای
AttributeErrorگرفتهنشده از بین برود، زیرا target باید در یک ماژول قابل ایمپورت تعریف شده باشد تا در حین بازسازی از pickle بارگذاری شود.مثالی از این خطای غیرقابلگرفتن از فرزند:
>>> import multiprocessing as mp >>> def knigit(): ... print("Ni!") ... >>> process = mp.Process(target=knigit) >>> process.start() >>> Traceback (most recent call last): File ".../multiprocessing/spawn.py", line ..., in spawn_main File ".../multiprocessing/spawn.py", line ..., in _main AttributeError: module '__main__' has no attribute 'knigit' >>> process <SpawnProcess name='SpawnProcess-1' pid=379473 parent=378707 stopped exitcode=1>
روشهای شروع spawn و forkserver را ببینید. اگرچه در صورت استفاده از روش شروع
"fork"این محدودیت برقرار نیست، اما از پایتون3.14این روش دیگر در هیچ سکویی پیشفرض نیست. زمینهها و متدهای شروع را ببینید. همچنین gh-132898 را ببینید.تغییر یافته در نسخهی 3.3: پارامتر daemon اضافه شد.
- run()¶
متدی که فعالیت فرایند را نشان میدهد.
شما میتوانید این متد را در یک زیرکلاس بازنویسی کنید. متد استاندارد
run()شیء فراخوانیپذیر ارسالشده به سازندهی شیء بهعنوان آرگومان target را، در صورت وجود، با آرگومانهای ترتیبی و کلیدواژهای که بهترتیب از آرگومانهای args و kwargs گرفته میشوند، فراخوانی میکند.استفاده از یک فهرست یا تاپل بهعنوان آرگومان args ارسالشده به
Process، همان اثر را دارد.مثال:
>>> from multiprocessing import Process >>> p = Process(target=print, args=[1]) >>> p.run() 1 >>> p = Process(target=print, args=(1,)) >>> p.run() 1
- start()¶
فعالیت فرایند را آغاز میکند.
این باید حداکثر یکبار برای هر شیء فرایند فراخوانی شود. این متد ترتیبی میدهد تا متد
run()شیء در یک فرایند جداگانه فراخوانی شود.
- join([timeout])¶
اگر آرگومان اختیاری timeout برابر
Noneباشد (مقدار پیشفرض)، این متد تا خاتمه یافتن فرایندی که متدjoin()آن فراخوانی میشود، مسدود میشود. اگر timeout یک عدد مثبت باشد، حداکثر به مدت timeout ثانیه مسدود میشود. توجه داشته باشید که این متد در صورت خاتمه یافتن فرایند یا به پایان رسیدن مهلت متد،Noneبرمیگرداند. برای تشخیص اینکه آیا فرایند خاتمه یافته است،exitcodeفرایند را بررسی کنید.میتوان یک فرایند را چندین بار join کرد.
یک فرایند نمیتواند به خودش بپیوندد، زیرا این کار باعث بنبست میشود. تلاش برای پیوستن به یک فرایند پیش از آنکه آغاز شده باشد، خطا است.
- name¶
نام فرایند. نام، رشتهای است که فقط برای اهداف شناسایی استفاده میشود. این نام هیچ معنایی ندارد. ممکن است به چندین فرایند نام یکسانی داده شود.
نام اولیه توسط سازنده تنظیم میشود. اگر نام صریحی به سازنده ارائه نشود، نامی به شکل 'Process-N1:N2:...:Nk' ساخته میشود، که در آن هر Nk، N-امین فرزند والد خود است.
- is_alive()¶
برمیگرداند که آیا فرایند زنده است یا خیر.
تقریباً، یک شیء فرایند از لحظهای که متد
start()برمیگردد تا زمانی که فرایند فرزند پایان مییابد، زنده است.
- daemon¶
پرچم daemon فرایند، یک مقدار بولی است. این پرچم باید پیش از فراخوانی
start()تنظیم شود.مقدار اولیه از فرآیند ایجادکننده به ارث میرسد.
هنگامی که یک فرایند خارج میشود، تلاش میکند تمام فرایندهای فرزند دیمونی (daemonic) خود را خاتمه دهد.
توجه داشته باشید که یک فرایند daemon (daemonic process) اجازه ندارد فرایندهای فرزند ایجاد کند. در غیر این صورت، اگر فرایند daemon هنگام خروج فرایند والد خود خاتمه یابد، فرایندهای فرزند خود را یتیم رها خواهد کرد. علاوه بر این، اینها daemonها یا سرویسهای Unix نیستند، بلکه فرایندهای عادی هستند که اگر فرایندهای غیر daemon (non-daemonic processes) خارج شده باشند، خاتمه داده خواهند شد (و join نخواهند شد).
علاوه بر API مربوط به
threading.Thread، اشیایProcessنیز از ویژگیها و متدهای زیر پشتیبانی میکنند:- pid¶
شناسه فرایند را برمیگرداند. پیش از ایجاد فرایند، این مقدار
Noneخواهد بود.
- exitcode¶
کد خروجی فرزند. اگر فرایند هنوز به پایان نرسیده باشد، این مقدار
Noneخواهد بود.اگر متد
run()فرزند بهطور عادی بازگشت کرد، کد خروجی ۰ خواهد بود. اگر از طریقsys.exit()با یک آرگومان عدد صحیح N خاتمه یافت، کد خروجی N خواهد بود.اگر فرزند بهدلیل استثنایی که در
run()گرفته نشده باشد خاتمه یابد، کد خروجی ۱ خواهد بود. اگر با سیگنال N خاتمه یافته باشد، کد خروجی مقدار منفی -N خواهد بود.
- authkey¶
کلید احراز هویت فرایند (یک رشته بایتی).
هنگامی که
multiprocessingمقداردهی اولیه میشود، یک رشته تصادفی با استفاده ازos.urandom()به فرایند اصلی اختصاص داده میشود.هنگامی که یک شیء
Processایجاد میشود، کلید احراز هویت فرایند والد خود را به ارث میبرد، اگرچه میتوان آن را با تنظیمauthkeyبه یک رشته بایتی دیگر تغییر داد.کلیدهای احراز هویت را ببینید.
- sentinel¶
یک دسته عددی از یک شیء سیستمی که هنگام پایان یافتن فرآیند «آماده» خواهد شد.
اگر میخواهید با استفاده از
multiprocessing.connection.wait()همزمان منتظر چند رویداد بمانید، میتوانید از این مقدار استفاده کنید. در غیر این صورت، فراخوانیjoin()سادهتر است.در ویندوز، این یک دستهی سیستمعامل است که با خانوادهی فراخوانیهای API
WaitForSingleObjectوWaitForMultipleObjectsقابل استفاده است. در POSIX، این یک توصیفگر پرونده است که با اولیههای ماژولselectقابل استفاده است.اضافه شده در نسخهی 3.3.
- interrupt()¶
فرایند را خاتمه میدهد. در POSIX با استفاده از سیگنال
SIGINTکار میکند. رفتار در ویندوز تعریفنشده است.بهطور پیشفرض، این کار فرایند فرزند را با پرتاب
KeyboardInterruptخاتمه میدهد. این رفتار را میتوان با تنظیم هندلر سیگنال مربوطه در فرایند فرزند باsignal.signal()برایSIGINTتغییر داد.توجه: اگر فرایند فرزند
KeyboardInterruptرا بگیرد و نادیده بگیرد، فرایند خاتمه نخواهد یافت.نکته: رفتار پیشفرض همچنین
exitcodeرا روی1تنظیم میکند، گویی یک استثنای گرفتهنشده در فرایند فرزند پرتاب شده است. برای داشتنexitcodeمتفاوت، میتوانید بهسادگیKeyboardInterruptرا بگیرید وexit(your_code)را فراخوانی کنید.اضافه شده در نسخهی 3.14.
- terminate()¶
فرایند را خاتمه میدهد. در POSIX این کار با استفاده از سیگنال
SIGTERMانجام میشود؛ در ویندوز ازTerminateProcess()استفاده میشود. توجه داشته باشید که هندلرهای خروج و بندهای finally و غیره اجرا نخواهند شد.توجه داشته باشید که فرایندهای زیرمجموعهی آن فرایند خاتمه داده نخواهند شد — بلکه صرفاً یتیم خواهند شد.
هشدار
اگر این متد زمانی استفاده شود که فرایند مرتبط از یک پایپ یا صف استفاده میکند، ممکن است آن پایپ یا صف خراب شود و برای فرایندهای دیگر غیرقابل استفاده گردد. بهطور مشابه، اگر فرایند یک قفل یا سمافور و غیره را کسب کرده باشد، خاتمه دادن به آن ممکن است باعث شود فرایندهای دیگر دچار بنبست شوند.
- kill()¶
مانند
terminate()، اما با استفاده از سیگنالSIGKILLدر POSIX.اضافه شده در نسخهی 3.7.
- close()¶
شیء
Processرا ببندید و تمام منابع مرتبط با آن را آزاد کنید. اگر فرایند زیربنایی هنوز در حال اجرا باشد،ValueErrorپرتاب میشود. پس از اینکهclose()با موفقیت برمیگردد، بیشتر متدها و ویژگیهای دیگر شیءProcess،ValueErrorرا پرتاب خواهند کرد.اضافه شده در نسخهی 3.7.
توجه داشته باشید که متدهای
start()،join()،is_alive()،terminate()وexitcodeفقط باید توسط فرایندی فراخوانی شوند که شیء فرایند را ایجاد کرده است.نمونهای از کاربرد برخی از متدهای
Process:>>> import multiprocessing, time, signal >>> mp_context = multiprocessing.get_context('spawn') >>> p = mp_context.Process(target=time.sleep, args=(1000,)) >>> print(p, p.is_alive()) <...Process ... initial> False >>> p.start() >>> print(p, p.is_alive()) <...Process ... started> True >>> p.terminate() >>> time.sleep(0.1) >>> print(p, p.is_alive()) <...Process ... stopped exitcode=-SIGTERM> False >>> p.exitcode == -signal.SIGTERM True
- exception multiprocessing.ProcessError¶
کلاس پایهی تمام استثناهای
multiprocessing.
- exception multiprocessing.BufferTooShort¶
استثنایی که توسط
Connection.recv_bytes_into()پرتاب میشود؛ زمانی که شیء بافر ارائهشده برای پیامی که خوانده میشود، بیش از حد کوچک باشد.اگر
eنمونهای ازBufferTooShortباشد، آنگاهe.args[0]پیام را بهصورت یک رشته بایتی برمیگرداند.
- exception multiprocessing.AuthenticationError¶
در صورت بروز خطای احراز هویت، پرتاب میشود.
- exception multiprocessing.TimeoutError¶
هنگامی که مهلت زمانی به پایان برسد، توسط متدهای دارای مهلت زمانی پرتاب میشود.
پایپها و صفها¶
هنگام استفاده از چندین فرایند، معمولاً برای ارتباط بین فرایندها از ارسال پیام استفاده میشود و از بهکارگیری هرگونه سازوکار اولیه همگامسازی مانند قفلها پرهیز میشود.
برای انتقال پیامها میتوانید از Pipe() (برای ارتباط بین دو فرآیند) یا یک صف (که امکان استفاده چندین تولیدکننده و مصرفکننده را فراهم میکند) استفاده کنید.
انواع Queue، SimpleQueue و JoinableQueue صفهای FIFO با چند تولیدکننده و چند مصرفکننده هستند که بر اساس کلاس queue.Queue در کتابخانه استاندارد ساخته شدهاند. تفاوت آنها در این است که Queue فاقد متدهای task_done() و join() است که در پایتون 2.5 به کلاس queue.Queue افزوده شدند.
اگر از JoinableQueue استفاده میکنید، باید به ازای هر وظیفهای که از صف برداشته میشود، JoinableQueue.task_done() را فراخوانی کنید؛ در غیر این صورت، سمافور استفادهشده برای شمارش تعداد وظایف تمامنشده ممکن است در نهایت سرریز کند و باعث پرتاب یک استثنا شود.
یکی از تفاوتها با سایر پیادهسازیهای صف در پایتون این است که صفهای multiprocessing تمام اشیایی را که در آنها قرار داده میشوند، با استفاده از pickle سریالسازی میکنند. شیءای که متد get بازمیگرداند، یک شیء بازسازیشده است که هیچ حافظهای را با شیء اصلی به اشتراک نمیگذارد.
توجه داشته باشید که میتوانید با استفاده از یک شیء مدیر، یک صف مشترک نیز ایجاد کنید -- به مدیرها مراجعه کنید.
توجه
multiprocessing از استثناهای معمول queue.Empty و queue.Full برای اعلام مهلت زمانی استفاده میکند. این استثناها در فضای نام multiprocessing در دسترس نیستند، بنابراین باید آنها را از queue ایمپورت کنید.
توجه
هنگامی که یک شیء در یک صف قرار میگیرد، آن شیء پیکل میشود و یک نخ پسزمینه بعداً دادههای پیکلشده را به یک پایپ زیرساختی تخلیه میکند. این موضوع پیامدهایی دارد که کمی تعجبآور هستند، اما نباید هیچ دشواری عملیای ایجاد کنند — اگر واقعاً شما را آزار میدهند، میتوانید در عوض از صفی استفاده کنید که با یک مدیر ایجاد شده است.
پس از قرار دادن یک شیء در یک صف خالی، ممکن است تأخیر بسیار ناچیزی پیش از آنکه متد
empty()صف مقدارFalseرا برگرداند وget_nowait()بتواند بدون اینکهqueue.Emptyپرتاب شود برگردد، وجود داشته باشد.اگر چندین فرایند شیءها را در صف قرار دهند، ممکن است شیءها در سمت دیگر خارج از ترتیب دریافت شوند. با این حال، اشیایی که توسط یک فرایند یکسان در صف قرار داده شدهاند، همیشه نسبت به یکدیگر در ترتیب مورد انتظار خواهند بود.
هشدار
اگر یک فرایند در حالی که سعی میکند از یک Queue استفاده کند، با استفاده از Process.terminate() یا os.kill() خاتمه داده شود، دادههای درون صف احتمالاً خراب میشوند. این ممکن است باعث شود هر فرایند دیگری هنگامی که بعداً سعی میکند از صف استفاده کند، با یک استثنا مواجه شود.
هشدار
همانطور که در بالا ذکر شد، اگر یک فرایند فرزند آیتمهایی را در یک صف قرار داده باشد (و از JoinableQueue.cancel_join_thread استفاده نکرده باشد)، آن فرایند تا زمانی که تمام آیتمهای بافرشده به پایپ تخلیه نشده باشند، خاتمه نخواهد یافت.
این بدان معناست که اگر سعی کنید آن فرایند را join کنید، ممکن است با بنبست مواجه شوید، مگر اینکه مطمئن باشید تمام آیتمهایی که در صف قرار داده شدهاند، مصرف شدهاند. بهطور مشابه، اگر فرایند فرزند غیردیمونی (non-daemonic) باشد، فرایند والد ممکن است هنگام خروج، هنگامی که سعی میکند تمام فرزندان غیردیمونی خود را join کند، معلق بماند.
توجه داشته باشید که یک صف ایجادشده با استفاده از مدیر، این مشکل را ندارد. دستورالعملهای برنامهنویسی را ببینید.
برای دیدن نمونهای از استفاده از صفها برای ارتباط بینفرایندی، مثالها را ببینید.
- multiprocessing.Pipe(duplex=True)¶
یک جفت
(conn1, conn2)از اشیایConnectionرا برمیگرداند که نشاندهندهی دو سر یک پایپ هستند.اگر duplex برابر
Trueباشد (پیشفرض)، پایپ دوطرفه است. اگر duplex برابرFalseباشد، پایپ یکطرفه است: ازconn1فقط برای دریافت پیامها و ازconn2فقط برای ارسال پیامها میتوان استفاده کرد.متد
send()شیء را با استفاده ازpickleسریالسازی میکند وrecv()شیء را بازسازی میکند.
- class multiprocessing.Queue([maxsize])¶
یک صف مشترک بین فرایندها را برمیگرداند که با استفاده از یک پایپ و چند قفل/سمافور پیادهسازی شده است. هنگامی که یک فرایند برای نخستین بار یک آیتم را در صف قرار میدهد، یک نخ تغذیهکننده آغاز میشود که شیءها را از یک بافر به درون پایپ منتقل میکند.
نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
استثناهای معمول
queue.Emptyوqueue.Fullاز ماژولqueueدر کتابخانهی استاندارد، برای اعلام مهلتهای زمانی پرتاب میشوند.Queueتمام متدهایqueue.Queueرا پیادهسازی میکند، بهجزtask_done()،join()وshutdown().- qsize()¶
اندازه تقریبی صف را برمیگرداند. به دلیل معناشناسی چندنخی/چندفرایندی، این عدد قابل اتکا نیست.
توجه داشته باشید که این ممکن است در سکوهایی مانند macOS، که
sem_getvalue()در آنها پیادهسازی نشده است،NotImplementedErrorرا پرتاب کند.
- empty()¶
اگر صف خالی باشد،
Trueو در غیر این صورتFalseرا برمیگرداند. به دلیل معناشناسی چندریسمانی/چندفرآیندی، این قابل اتکا نیست.ممکن است در صفهای بسته، یک
OSErrorپرتاب کند. (تضمین نمیشود)
- full()¶
اگر صف پر باشد،
Trueرا برمیگرداند، در غیر این صورتFalseرا برمیگرداند. به دلیل معناشناسی چندریسمانی/چندفرآیندی، این قابل اتکا نیست.
- put(obj[, block[, timeout]])¶
obj را در صف قرار میدهد. اگر آرگومان اختیاری block برابر
True(پیشفرض) و timeout برابرNone(پیشفرض) باشد، در صورت لزوم تا زمانی که یک جای خالی در دسترس باشد، مسدود میشود. اگر timeout یک عدد مثبت باشد، حداکثر به مدت timeout ثانیه مسدود میشود و اگر در آن مدت هیچ جای خالیای در دسترس نباشد، استثنایqueue.Fullرا پرتاب میکند. در غیر این صورت (block برابرFalseاست)، اگر یک جای خالی بلافاصله در دسترس باشد، آیتمی را در صف قرار میدهد، وگرنه استثنایqueue.Fullرا پرتاب میکند (در این حالت timeout نادیده گرفته میشود).تغییر یافته در نسخهی 3.8: اگر صف بسته شده باشد، بهجای
AssertionError،ValueErrorپرتاب میشود.
- put_nowait(obj)¶
معادل
put(obj, False)است.
- get([block[, timeout]])¶
یک آیتم را از صف حذف کرده و برمیگرداند. اگر آرگومان اختیاری block برابر
True(پیشفرض) و timeout برابرNone(پیشفرض) باشد، در صورت لزوم تا زمانی که آیتمی در دسترس باشد مسدود میشود. اگر timeout یک عدد مثبت باشد، حداکثر به مدت timeout ثانیه مسدود میشود و اگر هیچ آیتمی در آن مدت در دسترس نباشد، استثنایqueue.Emptyرا پرتاب میکند. در غیر این صورت (block برابرFalseاست)، اگر آیتمی بلافاصله در دسترس باشد، یک آیتم را برمیگرداند، در غیر این صورت استثنایqueue.Emptyرا پرتاب میکند (در این حالت timeout نادیده گرفته میشود).تغییر یافته در نسخهی 3.8: اگر صف بسته شده باشد، بهجای
OSError،ValueErrorپرتاب میشود.
- get_nowait()¶
معادل
get(False).
multiprocessing.Queueدارای چند متد اضافی است که درqueue.Queueوجود ندارند. این متدها معمولاً در بیشتر کدها غیرضروری هستند:- close()¶
بستن صف: آزاد کردن منابع داخلی.
پس از بستهشدن یک صف، دیگر نباید از آن استفاده شود. برای مثال، متدهای
get()،put()وempty()دیگر نباید فراخوانی شوند.نخ پسزمینه پس از آنکه تمام دادههای بافرشده را به پایپ تخلیه کرد، خارج میشود. این عملیات بهطور خودکار هنگام زبالهروبی صف فراخوانی میشود.
- join_thread()¶
به نخ پسزمینه بپیوندید. این کار فقط پس از فراخوانی
close()قابل انجام است. این عملیات تا خروج نخ پسزمینه مسدود میشود و اطمینان حاصل میکند که همهی دادههای موجود در بافر در پایپ نوشته شدهاند.بهطور پیشفرض، اگر یک فرایند ایجادکنندهی صف نباشد، در هنگام خروج تلاش میکند به نخ پسزمینهی صف ملحق شود. این فرایند میتواند
cancel_join_thread()را فراخوانی کند تاjoin_thread()هیچ کاری انجام ندهد.
- cancel_join_thread()¶
از مسدود شدن
join_thread()جلوگیری میکند. بهویژه، این امر مانع از پیوستن خودکار نخ پسزمینه هنگام خروج فرایند میشود —join_thread()را ببینید.ممکن است نام بهتری برای این متد
allow_exit_without_flush()باشد. احتمال دارد این متد باعث از دست رفتن دادههای در صف شود و شما تقریباً بهطور قطع نیازی به استفاده از آن نخواهید داشت. این متد در واقع فقط برای حالتی وجود دارد که لازم باشد فرآیند جاری بلافاصله و بدون انتظار برای تخلیهی دادههای در صف به پایپ زیرین خارج شود، و دادههای از دست رفته برایتان اهمیتی ندارد.
توجه
عملکرد این کلاس به یک پیادهسازی سمافور اشتراکیِ کارا در سیستمعامل میزبان نیاز دارد. بدون آن، عملکرد این کلاس غیرفعال خواهد شد و تلاش برای نمونهسازی از
QueueبهImportErrorمنجر خواهد شد. برای اطلاعات بیشتر bpo-3770 را ببینید. همین موضوع برای هر یک از انواع صف تخصصی فهرستشده در زیر نیز صادق است.
- class multiprocessing.SimpleQueue¶
این یک نوع سادهشده از
Queueاست، بسیار شبیه به یکPipeقفلشده.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
- close()¶
بستن صف: آزاد کردن منابع داخلی.
پس از بستهشدن یک صف، دیگر نباید از آن استفاده شود. برای مثال، متدهای
get()،put()وempty()نباید دیگر فراخوانی شوند.اضافه شده در نسخهی 3.9.
- empty()¶
اگر صف خالی باشد،
Trueو در غیر این صورتFalseرا برمیگرداند.همیشه در صورت بسته بودن SimpleQueue، یک
OSErrorپرتاب میکند.
- get()¶
یک آیتم را از صف حذف میکند و برمیگرداند.
- put(item)¶
آیتم را در صف قرار دهید.
- class multiprocessing.JoinableQueue([maxsize])¶
JoinableQueue، زیرکلاسی ازQueue، صفی است که علاوه بر این دارای متدهایtask_done()وjoin()است.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
- task_done()¶
نشان میدهد که وظیفهای که پیشتر در صف قرار گرفته است، کامل شده است. توسط مصرفکنندگان صف استفاده میشود. برای هر
get()که برای واکشی یک وظیفه استفاده میشود، فراخوانی بعدیtask_done()به صف میگوید که پردازش روی آن وظیفه کامل شده است.اگر یک
join()در حال حاضر مسدود باشد، زمانی ادامه مییابد که همهی آیتمها پردازش شده باشند (یعنی برای هر آیتمی که باput()در صف قرار گرفته است، یک فراخوانیtask_done()دریافت شده باشد).اگر بیش از تعداد آیتمهای قرار دادهشده در صف فراخوانی شود، یک
ValueErrorپرتاب میکند.
- join()¶
مسدود میشود تا تمام آیتمهای صف دریافت و پردازش شده باشند.
تعداد کارهای ناتمام هرگاه که یک آیتم به صف اضافه شود، افزایش مییابد. هرگاه یک مصرفکننده
task_done()را فراخوانی کند تا نشان دهد آن آیتم دریافت شد و همهی کارهای مربوط به آن کامل شده است، این تعداد کاهش مییابد. هنگامی که تعداد کارهای ناتمام به صفر برسد،join()از حالت مسدود خارج میشود.
متفرقه¶
- multiprocessing.active_children()¶
فهرستی از تمام فرزندان زندهی فرایند جاری را برمیگرداند.
فراخوانی این، دارای اثر جانبیِ «پیوستن» به هر یک از فرآیندهایی است که از قبل به پایان رسیدهاند.
- multiprocessing.cpu_count()¶
تعداد پردازندههای سیستم را برمیگرداند.
این عدد معادل تعداد پردازندههایی که فرایند جاری میتواند از آنها استفاده کند نیست. تعداد پردازندههای قابل استفاده را میتوان با
os.process_cpu_count()(یاlen(os.sched_getaffinity(0))) به دست آورد.هنگامی که نتوان تعداد پردازندهها را تعیین کرد، یک
NotImplementedErrorپرتاب میشود.همچنین ملاحظه نمائید
تغییر یافته در نسخهی 3.13: مقدار بازگشتی همچنین میتواند با استفاده از پرچم
-X cpu_countیاPYTHON_CPU_COUNTبازنویسی شود، زیرا این صرفاً پوششی بر APIهای شمارش پردازنده درosاست.
- multiprocessing.current_process()¶
شیء
Processمتناظر با فرایند جاری را برمیگرداند.معادلی از
threading.current_thread().
- multiprocessing.parent_process()¶
شیء
Processمتناظر با فرایند والدِcurrent_process()را برمیگرداند. برای فرایند اصلی،parent_processبرابرNoneخواهد بود.اضافه شده در نسخهی 3.8.
- multiprocessing.freeze_support()¶
افزودن پشتیبانی برای زمانی که برنامهای که از
multiprocessingاستفاده میکند، فریزشده باشد تا یک پرونده اجرایی تولید کند. (با py2exe، PyInstaller و cx_Freeze آزمایش شده است.)باید این تابع را بلافاصله پس از خط
if __name__ == '__main__'در ماژول اصلی فراخوانی کنید. برای مثال:from multiprocessing import Process, freeze_support def f(): print('hello world!') if __name__ == '__main__': freeze_support() Process(target=f).start()
اگر ردیف
freeze_support()حذف شود، تلاش برای اجرای پرونده اجرایی فریزشده موجب پرتابRuntimeErrorمیشود.فراخوانی
freeze_support()هنگامی که روش راهاندازی spawn نباشد، هیچ تأثیری ندارد. علاوه بر این، اگر ماژول بهصورت عادی توسط مفسر پایتون اجرا شود (برنامه فریز نشده باشد)،freeze_support()هیچ تأثیری ندارد.
- multiprocessing.get_all_start_methods()¶
فهرستی از روشهای شروع پشتیبانیشده را برمیگرداند که اولین آنها پیشفرض است. روشهای شروع ممکن
'fork'،'spawn'و'forkserver'هستند. همهی سکوها از همهی روشها پشتیبانی نمیکنند. به زمینهها و متدهای شروع مراجعه کنید.اضافه شده در نسخهی 3.4.
- multiprocessing.get_context(method=None)¶
یک شیء زمینه برمیگرداند که همان ویژگیهای ماژول
multiprocessingرا دارد.اگر method برابر
Noneباشد، زمینه پیشفرض برگردانده میشود. توجه داشته باشید که اگر متد شروع سراسری تنظیم نشده باشد، این کار آن را روی پیشفرض سیستم تنظیم میکند. برای جزئیات بیشتر، روش شروع سراسری را ببینید. در غیر این صورت، method باید'fork'،'spawn'یا'forkserver'باشد. اگر متد شروع مشخصشده در دسترس نباشد،ValueErrorپرتاب میشود. زمینهها و متدهای شروع را ببینید.اضافه شده در نسخهی 3.4.
- multiprocessing.get_start_method(allow_none=False)¶
نام روش شروع استفادهشده برای شروع فرآیندها را برمیگرداند.
اگر روش شروع سراسری تنظیم نشده باشد و allow_none برابر
Falseباشد، روش شروع سراسری روی مقدار پیشفرض تنظیم میشود و نام آن برگردانده میشود. برای جزئیات بیشتر روش شروع سراسری را ببینید.مقدار بازگشتی میتواند
'fork'،'spawn'،'forkserver'یاNoneباشد. به زمینهها و متدهای شروع مراجعه کنید.اضافه شده در نسخهی 3.4.
تغییر یافته در نسخهی 3.8: در macOS، روش راهاندازی spawn اکنون پیشفرض است. روش راهاندازی fork باید ناامن در نظر گرفته شود، زیرا میتواند منجر به فروپاشی زیرفرایند شود. bpo-33725 را ببینید.
- multiprocessing.set_executable(executable)¶
مسیر مفسر پایتون را برای استفاده هنگام راهاندازی یک فرایند فرزند تنظیم کنید. (بهطور پیشفرض از
sys.executableاستفاده میشود). تعبیهکنندهها (embedders) احتمالاً باید کاری مانند زیر انجام دهندset_executable(os.path.join(sys.exec_prefix, 'pythonw.exe'))
پیش از آنکه بتوانند فرآیندهای فرزند ایجاد کنند.
تغییر یافته در نسخهی 3.4: اکنون در POSIX در صورتی پشتیبانی میشود که از روش راهاندازی
'spawn'استفاده شود.تغییر یافته در نسخهی 3.11: یک شیء شبهمسیر را میپذیرد.
- multiprocessing.set_forkserver_preload(module_names)¶
فهرستی از نام ماژولها را برای فرایند اصلی forkserver تنظیم کنید تا تلاش کند آنها را ایمپورت کند، بهطوری که وضعیت از پیش ایمپورتشدهی آنها توسط فرایندهای forkشده به ارث برده شود. هرگونه
ImportErrorهنگام انجام این کار بهصورت بیسر و صدا نادیده گرفته میشود. از این میتوان بهعنوان بهبود کارایی برای اجتناب از کار تکراری در هر فرایند استفاده کرد.برای این که این کار کند، باید پیش از راهاندازی فرایند forkserver فراخوانی شود (پیش از ایجاد یک
Poolیا آغاز یکProcess).تنها زمانی معنا دارد که از متد شروع
'forkserver'استفاده شود. زمینهها و متدهای شروع را ببینید.اضافه شده در نسخهی 3.4.
- multiprocessing.set_start_method(method, force=False)¶
متدی را که باید برای شروع فرآیندهای فرزند استفاده شود، تنظیم کنید. آرگومان method میتواند
'fork'،'spawn'یا'forkserver'باشد. اگر متد شروع از قبل تنظیمشده باشد و force برابرTrueنباشد،RuntimeErrorپرتاب میشود. اگر method برابرNoneو force برابرTrueباشد، متد شروع بهNoneتنظیم میشود. اگر method برابرNoneو force برابرFalseباشد، زمینه به زمینه پیشفرض تنظیم میشود.توجه داشته باشید که این باید حداکثر یک بار فراخوانی شود، و باید در بند
if __name__ == '__main__'ماژول اصلی محافظت شود.زمینهها و متدهای شروع را ببینید.
اضافه شده در نسخهی 3.4.
توجه
multiprocessing شامل هیچ معادلی برای threading.active_count()، threading.enumerate()، threading.settrace()، threading.setprofile()، threading.Timer یا threading.local نیست.
اشیای اتصال¶
اشیاء Connection امکان ارسال و دریافت اشیاء پیکلپذیر یا رشتهها را فراهم میکنند. میتوان آنها را بهعنوان سوکتهای متصل پیاممحور در نظر گرفت.
اشیای Connection معمولاً با استفاده از Pipe ایجاد میشوند -- همچنین شنوندهها و کلاینتها را ببینید.
- class multiprocessing.connection.Connection¶
- send(obj)¶
یک شیء را به سمت دیگر اتصال ارسال کنید که باید با
recv()خوانده شود.شیء باید پیکلپذیر باشد. پیکلهای بسیار بزرگ (تقریباً ۳۲ MiB یا بیشتر، هرچند بستگی به سیستمعامل دارد) ممکن است باعث پرتاب استثنای
ValueErrorشوند.
- recv()¶
شیء ارسالشده از طرف دیگر اتصال با استفاده از
send()را بازمیگرداند. تا زمانی که چیزی برای دریافت وجود نداشته باشد، مسدود میماند. اگر هیچ چیزی برای دریافت باقی نمانده باشد و طرف دیگر اتصال بسته شده باشد،EOFErrorرا پرتاب میکند.
- fileno()¶
توصیفگر پرونده یا دسته مورد استفاده توسط اتصال را برمیگرداند.
- close()¶
اتصال را ببندید.
این بهطور خودکار هنگام زبالهروبی اتصال فراخوانی میشود.
- poll([timeout])¶
برمیگرداند که آیا دادهای برای خواندن در دسترس است یا خیر.
اگر timeout مشخص نشده باشد، بلافاصله بازمیگردد. اگر timeout یک عدد باشد، حداکثر مهلت زمانی برای مسدود شدن بر حسب ثانیه را مشخص میکند. اگر timeout برابر
Noneباشد، از یک مهلت زمانی نامحدود استفاده میشود.توجه داشته باشید که با استفاده از
multiprocessing.connection.wait()میتوان چندین شیء اتصال را بهطور همزمان پایش کرد.
- send_bytes(buf[, offset[, size]])¶
داده بایتی را از یک bytes-like object بهعنوان یک پیام کامل ارسال کنید.
اگر offset داده شود، داده از آن موقعیت در buf خوانده میشود. اگر size داده شود، همان تعداد بایت از buf خوانده خواهد شد. بافرهای بسیار بزرگ (حدود ۳۲ MiB+، هرچند به سیستمعامل بستگی دارد) ممکن است یک استثنا
ValueErrorرا پرتاب کنند
- recv_bytes([maxlength])¶
یک پیام کامل از دادههای بایتی ارسالشده از سمت دیگر اتصال را بهعنوان یک رشته برمیگرداند. تا زمانی که چیزی برای دریافت وجود نداشته باشد، مسدود میشود. اگر دیگر چیزی برای دریافت باقی نمانده باشد و سمت دیگر بسته شده باشد،
EOFErrorرا پرتاب میکند.اگر maxlength تعیین شده باشد و پیام طولانیتر از maxlength باشد،
OSErrorپرتاب میشود و اتصال دیگر قابل خواندن نخواهد بود.
- recv_bytes_into(buf[, offset])¶
یک پیام کامل از دادههای بایتی ارسالشده از طرف دیگر اتصال را در buf میخواند و تعداد بایتهای پیام را برمیگرداند. تا زمانی که چیزی برای دریافت وجود داشته باشد، مسدود میشود. اگر دیگر چیزی برای دریافت باقی نمانده باشد و طرف دیگر اتصال بسته شده باشد،
EOFErrorرا پرتاب میکند.buf باید یک bytes-like object قابلنوشتن باشد. اگر offset داده شود، پیام از آن موقعیت در بافر نوشته خواهد شد. آفست باید یک عدد صحیح نامنفی باشد که از طول buf (بر حسب بایت) کمتر باشد.
اگر بافربیش از حد کوتاه باشد، استثنای
BufferTooShortپرتاب میشود و پیام کامل بهصورتe.args[0]در دسترس است، که در آنeنمونهی استثنا است.
تغییر یافته در نسخهی 3.3: اکنون میتوان خود اشیای Connection را با استفاده از
Connection.send()وConnection.recv()بین فرایندها منتقل کرد.اشیای اتصال اکنون از پروتکل مدیریت زمینه نیز پشتیبانی میکنند — Context Manager Types را ببینید.
__enter__()شیء اتصال را برمیگرداند و__exit__()close()را فراخوانی میکند.
برای مثال:
>>> from multiprocessing import Pipe
>>> a, b = Pipe()
>>> a.send([1, 'hello', None])
>>> b.recv()
[1, 'hello', None]
>>> b.send_bytes(b'thank you')
>>> a.recv_bytes()
b'thank you'
>>> import array
>>> arr1 = array.array('i', range(5))
>>> arr2 = array.array('i', [0] * 10)
>>> a.send_bytes(arr1)
>>> count = b.recv_bytes_into(arr2)
>>> assert count == len(arr1) * arr1.itemsize
>>> arr2
array('i', [0, 1, 2, 3, 4, 0, 0, 0, 0, 0])
هشدار
متد Connection.recv() بهطور خودکار دادههای دریافتی را از پیکل خارج میکند (unpickles)، که میتواند یک خطر امنیتی باشد، مگر اینکه بتوانید به فرایندی که پیام را ارسال کرده است اعتماد کنید.
بنابراین، مگر اینکه شیء اتصال با استفاده از Pipe() تولید شده باشد، باید تنها پس از انجام نوعی احراز هویت از متدهای recv() و send() استفاده کنید. کلیدهای احراز هویت را ببینید.
هشدار
اگر فرایندی در حالی که تلاش میکند از یک پایپبخواند یا در آن بنویسد خاتمه داده شود، احتمال دارد دادههای درون پایپ خراب شوند، زیرا ممکن است اطمینان از محل مرزهای پیام غیرممکن شود.
اولیههای همگامسازی¶
بهطور کلی، سازوکارهای همگامسازی در برنامههای چندفرآیندی به اندازه برنامههای چندنخی ضروری نیستند. مستندات ماژول threading را ببینید.
توجه داشته باشید که میتوان با استفاده از یک شیء مدیر، سازوکارهای همگامسازی را نیز ایجاد کرد -- به مدیرها مراجعه کنید.
- class multiprocessing.Barrier(parties[, action[, timeout]])¶
یک شیء سد (barrier): نسخهای از
threading.Barrier.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
اضافه شده در نسخهی 3.3.
- class multiprocessing.BoundedSemaphore([value])¶
یک شیء سمافور محدود: معادلی نزدیک به
threading.BoundedSemaphore.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
تنها یک تفاوت با مشابه نزدیک آن وجود دارد: نام نخستین آرگومان متد
acquireآن block است، همانگونه که باLock.acquire()سازگار است.- locked()¶
یک بولی برمیگرداند که نشان میدهد آیا این شیء هماکنون قفل است یا خیر.
اضافه شده در نسخهی 3.14.
توجه
در macOS، این از
Semaphoreقابل تشخیص نیست، زیراsem_getvalue()در آن پیادهسازی نشده است.
- class multiprocessing.Condition([lock])¶
یک متغیر شرطی: نام مستعاری برای
threading.Condition.اگر lock مشخص شده باشد، باید یک شیء
LockیاRLockازmultiprocessingباشد.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
تغییر یافته در نسخهی 3.3: متد
wait_for()افزوده شد.
- class multiprocessing.Event¶
یک نسخهی رونوشت (clone) از
threading.Event.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
- class multiprocessing.Lock¶
یک شیء قفل غیربازگشتی: مشابهی نزدیک به
threading.Lock. هنگامی که یک فرایند یا نخ یک قفل را کسب کرد، تلاشهای بعدی برای کسب آن از سوی هر فرایند یا نخ، تا زمانی که آزاد شود مسدود میشوند؛ هر فرایند یا نخ میتواند آن را آزاد کند. مفاهیم و رفتارهایthreading.Lock، همانگونه که در مورد نخها اعمال میشود، در اینجا درmultiprocessing.Lock، همانگونه که در مورد فرایندها یا نخها اعمال میشود، تکرار شدهاند، مگر آنگونه که ذکر شده باشد.توجه داشته باشید که
Lockدر واقع یک تابع کارخانهای است که نمونهای ازmultiprocessing.synchronize.Lockرا برمیگرداند که با یک زمینه پیشفرض مقداردهی اولیه شده است.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
Lockاز پروتکل context manager پشتیبانی میکند و بنابراین میتوان از آن در دستورهایwithاستفاده کرد.- acquire(block=True, timeout=None)¶
یک قفل را بهصورت مسدودکننده یا غیرمسدودکننده به دست آورید.
اگر آرگومان block روی
True(پیشفرض) تنظیم شده باشد، فراخوانی متد تا زمانی که قفل در وضعیت قفلنشده باشد مسدود میشود، سپس آن را در وضعیت قفلشده قرار میدهد وTrueرا برمیگرداند. توجه داشته باشید که نام نخستین آرگومان با نام آن درthreading.Lock.acquire()متفاوت است.وقتی آرگومان block روی
Falseتنظیم شده باشد، فراخوانی متد مسدود نمیشود. اگر قفل در حال حاضر در وضعیت قفلشده باشد،Falseرا برمیگرداند؛ در غیر این صورت قفل را در وضعیت قفلشده قرار میدهد وTrueرا برمیگرداند.هنگامی که با یک مقدار مثبت از نوع ممیز شناور برای timeout فراخوانی شود، تا زمانی که نمیتوان قفل را به دست آورد، حداکثر به تعداد ثانیههای مشخصشده توسط timeout مسدود میشود. فراخوانیهایی با مقدار منفی برای timeout معادل timeout برابر با ۰ هستند. فراخوانیهایی با مقدار
Noneبرای timeout (پیشفرض)، مهلت زمانی را روی بینهایت تنظیم میکنند. توجه داشته باشید که نحوهی برخورد با مقادیر منفی یاNoneبرای timeout با رفتار پیادهسازیشده درthreading.Lock.acquire()متفاوت است. اگر آرگومان block رویFalseتنظیم شده باشد، آرگومان timeout هیچ تأثیر عملی ندارد و بنابراین نادیده گرفته میشود. اگر قفل به دست آمده باشدTrueو اگر مهلت زمانی سپری شده باشدFalseبرمیگرداند.
- release()¶
آزاد کردن یک قفل. میتوان آن را از هر فرایند یا نخی فراخوانی کرد، نه فقط فرایند یا نخی که در ابتدا قفل را کسب کرده است.
رفتار همانند
threading.Lock.release()است، به جز اینکه هنگام فراخوانی روی یک قفل باز، یکValueErrorپرتاب میشود.
- locked()¶
یک بولی برمیگرداند که نشان میدهد آیا این شیء هماکنون قفل است یا خیر.
اضافه شده در نسخهی 3.14.
- class multiprocessing.RLock¶
یک شیء قفل بازگشتی: مشابهی نزدیک به
threading.RLock. یک قفل بازگشتی باید توسط فرایند یا نخی که آن را کسب کرده است آزاد شود. هنگامی که یک فرایند یا نخ یک قفل بازگشتی را کسب کرد، همان فرایند یا نخ میتواند دوباره آن را بدون مسدود شدن کسب کند؛ آن فرایند یا نخ باید به ازای هر بار که آن را کسب کرده است، یک بار آن را آزاد کند.توجه داشته باشید که
RLockدر واقع یک تابع کارخانه است که نمونهای ازmultiprocessing.synchronize.RLockرا برمیگرداند که با یک زمینه پیشفرض مقداردهی اولیه شده است.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
RLockاز پروتکل context manager پشتیبانی میکند، بنابراین میتوان از آن در دستورهایwithاستفاده کرد.- acquire(block=True, timeout=None)¶
یک قفل را بهصورت مسدودکننده یا غیرمسدودکننده به دست آورید.
هنگام فراخوانی با آرگومان block که روی
Trueتنظیمشده باشد، تا زمانی که قفل در وضعیت قفلنشده قرار گیرد (متعلق به هیچ فرایند یا نخی نباشد) مسدود میشود، مگر اینکه قفل از قبل متعلق به فرایند یا نخ جاری باشد. سپس فرایند یا نخ جاری مالکیت قفل را به دست میگیرد (اگر از قبل مالکیت آن را نداشته باشد) و سطح بازگشتی درون قفل یک واحد افزایش مییابد، که در نتیجه مقدار بازگشتیTrueخواهد بود. توجه داشته باشید که چندین تفاوت در رفتار این اولین آرگومان در مقایسه با پیادهسازیthreading.RLock.acquire()وجود دارد، که با نام خود آرگومان آغاز میشود.هنگامی که با آرگومان block که روی
Falseتنظیم شده باشد فراخوانی شود، مسدود نمیشود. اگر قفل از قبل توسط فرآیند یا نخ دیگری کسب شده باشد (و بنابراین تحت مالکیت آن باشد)، فرآیند یا نخ جاری مالکیت آن را به دست نمیآورد و سطح بازگشت درون قفل تغییر نمیکند و در نتیجه مقدار بازگشتیFalseخواهد بود. اگر قفل در وضعیت بازشده باشد، فرآیند یا نخ جاری مالکیت آن را به دست میآورد و سطح بازگشت افزایش مییابد و در نتیجه مقدار بازگشتیTrueخواهد بود.استفاده و رفتارهای آرگومان timeout همانند
Lock.acquire()است. توجه داشته باشید که برخی از این رفتارهای timeout با رفتارهای پیادهسازیشده درthreading.RLock.acquire()متفاوت است.
- release()¶
یک قفل را آزاد میکند و سطح بازگشت را کاهش میدهد. اگر پس از کاهش، سطح بازگشت صفر باشد، قفل را به حالت آزاد بازنشانی میکند (متعلق به هیچ فرایند یا نخی نیست) و اگر هر فرایند یا نخ دیگری برای آزاد شدن قفل مسدود شده و در انتظار باشد، دقیقاً به یکی از آنها اجازه میدهد ادامه دهد. اگر پس از کاهش، سطح بازگشت همچنان غیرصفر باشد، قفل همچنان قفلشده و متعلق به فرایند یا نخ فراخوان باقی میماند.
این متد را تنها زمانی فراخوانی کنید که فرایند یا نخ فراخوانیکننده، مالک قفل باشد. اگر این متد توسط فرایند یا نخی غیر از مالک فراخوانی شود یا قفل در وضعیت قفلنشده (بدون مالک) باشد، یک
AssertionErrorپرتاب میشود. توجه داشته باشید که نوع استثنای پرتابشده در این وضعیت با رفتار پیادهسازیشده درthreading.RLock.release()متفاوت است.
- locked()¶
یک بولی برمیگرداند که نشان میدهد آیا این شیء هماکنون قفل است یا خیر.
اضافه شده در نسخهی 3.14.
- class multiprocessing.Semaphore([value])¶
یک شیء سمافور: مشابهی نزدیک به
threading.Semaphore.نمونهسازی از این کلاس ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.
تنها یک تفاوت با مشابه نزدیک آن وجود دارد: نام نخستین آرگومان متد
acquireآن block است، همانگونه که باLock.acquire()سازگار است.- get_value()¶
مقدار فعلی سمافور را برمیگرداند.
توجه داشته باشید که این ممکن است در سکوهایی مانند macOS، که
sem_getvalue()در آنها پیادهسازی نشده است،NotImplementedErrorرا پرتاب کند.
- locked()¶
یک بولی برمیگرداند که نشان میدهد آیا این شیء هماکنون قفل است یا خیر.
اضافه شده در نسخهی 3.14.
توجه
در macOS، از sem_timedwait پشتیبانی نمیشود، بنابراین فراخوانی acquire() با یک مهلت زمانی، رفتار آن تابع را با استفاده از یک حلقهی توقف شبیهسازی خواهد کرد.
توجه
برخی از قابلیتهای این بسته به یک پیادهسازی سمافور اشتراکی نیاز دارند که روی سیستمعامل میزبان بهدرستی کار کند. بدون وجود آن، ماژول multiprocessing.synchronize غیرفعال خواهد شد و تلاشها برای ایمپورت آن منجر به ImportError خواهند شد. برای اطلاعات بیشتر bpo-3770 را ببینید.
مدیرها¶
مدیرها روشی برای ایجاد دادههایی فراهم میکنند که میتوان آنها را بین فرآیندهای مختلف به اشتراک گذاشت، از جمله اشتراکگذاری از طریق شبکه بین فرآیندهایی که روی ماشینهای مختلف اجرا میشوند. یک شیء مدیر یک فرآیند سرور را کنترل میکند که اشیاء مشترک را مدیریت میکند. سایر فرآیندها میتوانند با استفاده از پراکسیها به اشیاء مشترک دسترسی داشته باشند.
- multiprocessing.Manager()¶
یک شیء
SyncManagerراهاندازیشده برمیگرداند که میتوان از آن برای اشتراکگذاری اشیاء بین فرایندها استفاده کرد. شیء مدیر برگرداندهشده متناظر با یک فرایند فرزند ایجادشده است و متدهایی دارد که اشیاء مشترک را ایجاد میکنند و پراکسیهای متناظر را برمیگردانند.
فرایندهای مدیر به محض زبالهروبی شدن یا خارج شدن فرایند والدشان، خاموش میشوند. کلاسهای مدیر در ماژول multiprocessing.managers تعریف شدهاند:
- class multiprocessing.managers.BaseManager(address=None, authkey=None, serializer='pickle', ctx=None, *, shutdown_timeout=1.0)¶
یک شیء BaseManager ایجاد کنید.
پس از ایجاد، باید
start()یاget_server().serve_forever()را فراخوانی کرد تا اطمینان حاصل شود که شیء مدیر به یک فرایند مدیر آغازشده ارجاع دارد.address نشانیای است که فرایند مدیر برای اتصالهای جدید به آن گوش میدهد. اگر address برابر
Noneباشد، یک نشانی دلخواه انتخاب میشود.authkey کلید احراز هویتی است که برای بررسی اعتبار اتصالات ورودی به فرایند سرور استفاده میشود. اگر authkey برابر
Noneباشد، ازcurrent_process().authkeyاستفاده میشود. در غیر این صورت از authkey استفاده میشود و آن باید یک رشته بایتی باشد.serializer باید
'pickle'(استفاده از سریالسازیpickle) یا'xmlrpclib'(استفاده از سریالسازیxmlrpc.client) باشد.ctx یک شیء زمینه است، یا
None(از زمینه جاری استفاده میشود). اگرNoneباشد، این فراخوانی ممکن است متد شروع سراسری را تنظیم کند. برای جزئیات بیشتر، روش شروع سراسری را ببینید.shutdown_timeout یک مهلت زمانی بر حسب ثانیه است که در متد
shutdown()برای انتظار تا تکمیل شدن فرایند مورد استفاده مدیر به کار میرود. اگر shutdown با انقضای مهلت مواجه شود، فرایند خاتمه داده میشود. اگر خاتمه دادن فرایند نیز با انقضای مهلت مواجه شود، فرایند کشته میشود.تغییر یافته در نسخهی 3.11: پارامتر shutdown_timeout اضافه شد.
- start([initializer[, initargs]])¶
برای راهاندازی مدیر، یک زیرفرایند را آغاز کنید. اگر initializer برابر
Noneنباشد، زیرفرایند هنگام شروع،initializer(*initargs)را فراخوانی میکند.
- get_server()¶
یک شیء
Serverرا برمیگرداند که نشاندهنده سرور واقعی تحت کنترل Manager است. شیءServerاز متدserve_forever()پشتیبانی میکند:>>> from multiprocessing.managers import BaseManager >>> manager = BaseManager(address=('', 50000), authkey=b'abc') >>> server = manager.get_server() >>> server.serve_forever()
Serverهمچنین دارای یک ویژگیaddressاست.
- connect()¶
یک شیء مدیر محلی را به یک فرایند مدیر راه دور متصل کنید:
>>> from multiprocessing.managers import BaseManager >>> m = BaseManager(address=('127.0.0.1', 50000), authkey=b'abc') >>> m.connect()
- shutdown()¶
فرایند استفادهشده توسط مدیر را متوقف میکند. این تنها در صورتی در دسترس است که از
start()برای راهاندازی فرایند سرور استفاده شده باشد.این را میتوان چندین بار فراخوانی کرد.
- register(typeid[, callable[, proxytype[, exposed[, method_to_typeid[, create_method]]]]])¶
یک classmethod که میتوان از آن برای ثبت یک نوع یا فراخوانیپذیر در کلاس مدیر استفاده کرد.
typeid یک «شناسه نوع» است که برای شناسایی نوع خاصی از یک شیء مشترک استفاده میشود. این باید یک رشته باشد.
callable یک شیء فراخوانیپذیر است که برای ایجاد اشیاء برای این شناسهی نوع استفاده میشود. اگر یک نمونه مدیر قرار باشد با استفاده از متد
connect()به سرور متصل شود، یا اگر آرگومان create_method برابرFalseباشد، میتوان آن را بهصورتNoneباقی گذاشت.proxytype زیرکلاسی از
BaseProxyاست که برای ایجاد پراکسیهایی برای اشیاء اشتراکگذاریشده با این typeid استفاده میشود. اگرNoneباشد، یک کلاس پراکسی بهصورت خودکار ایجاد میشود.از exposed برای مشخص کردن دنبالهای از نام متدهایی استفاده میشود که پراکسیهای این typeid باید مجاز به دسترسی به آنها از طریق
BaseProxy._callmethod()باشند. (اگر exposed برابرNoneباشد، در صورت وجود، ازproxytype._exposed_بهجای آن استفاده میشود.) در صورتی که هیچ فهرست exposed مشخص نشده باشد، همه «متدهای عمومی» شیء مشترک قابل دسترسی خواهند بود. (در اینجا «متد عمومی» به هر ویژگی گفته میشود که دارای متد__call__()باشد و نام آن با'_'آغاز نشود.)method_to_typeid یک نگاشت است که برای مشخص کردن نوع بازگشتی آن دسته از متدهای آشکارشده که باید یک پراکسی بازگردانند، به کار میرود. این نگاشت، نام متدها را به رشتههای typeid نگاشت میکند. (اگر method_to_typeid برابر
Noneباشد، در صورت وجودproxytype._method_to_typeid_بهجای آن استفاده میشود.) اگر نام یک متد کلید این نگاشت نباشد یا نگاشتNoneباشد، شیء بازگرداندهشده توسط متد بهصورت مقدار کپی خواهد شد.create_method تعیین میکند که آیا باید متدی با نام typeid ایجاد شود که میتوان از آن برای درخواست از فرایند سرور جهت ایجاد یک شیء مشترک جدید و بازگرداندن یک پراکسی برای آن استفاده کرد. بهطور پیشفرض، این مقدار
Trueاست.
نمونههای
BaseManagerهمچنین دارای یک ویژگی فقطخواندنی هستند:- address¶
نشانی مورد استفادهی مدیر.
تغییر یافته در نسخهی 3.3: اشیاء مدیر از پروتکل مدیریت زمینه پشتیبانی میکنند -- Context Manager Types را ببینید.
__enter__()فرایند سرور را (اگر پیشتر آغاز نشده باشد) آغاز میکند و سپس شیء مدیر را برمیگرداند.__exit__()shutdown()را فراخوانی میکند.در نسخههای پیشین،
__enter__()فرایند سرور مدیر را در صورتی که از قبل آغاز نشده بود، آغاز نمیکرد.
- class multiprocessing.managers.SyncManager¶
زیرکلاسی از
BaseManagerکه میتوان از آن برای همگامسازی فرآیندها استفاده کرد. شیءهایی از این نوع توسطmultiprocessing.Manager()برگردانده میشوند.متدهای آن اشیای پراکسی را برای تعدادی از انواع دادهی رایج که باید بین فرآیندها همگامسازی شوند، ایجاد میکنند و بازمیگردانند. این بهویژه شامل فهرستها و دیکشنریهای مشترک میشود.
- Barrier(parties[, action[, timeout]])¶
یک شیء مشترک
threading.Barrierایجاد کنید و یک پراکسی برای آن برگردانید.اضافه شده در نسخهی 3.3.
- BoundedSemaphore([value])¶
یک شیء مشترک
threading.BoundedSemaphoreایجاد میکند و یک پراکسی برای آن برمیگرداند.
- Condition([lock])¶
یک شیء
threading.Conditionمشترک ایجاد کنید و یک پراکسی برای آن برگردانید.اگر lock ارائه شود، باید یک پراکسی برای یک شیء
threading.Lockیاthreading.RLockباشد.تغییر یافته در نسخهی 3.3: متد
wait_for()افزوده شد.
- Event()¶
یک شیء مشترک
threading.Eventایجاد میکند و یک پراکسی برای آن برمیگرداند.
- Lock()¶
یک شیء
threading.Lockمشترک ایجاد کنید و یک پراکسی برای آن برگردانید.
- Queue([maxsize])¶
یک شیء
queue.Queueمشترک ایجاد کنید و یک پراکسی برای آن برگردانید.
- RLock()¶
ایجاد یک شیء مشترک
threading.RLockو برگرداندن یک پراکسی برای آن.
- Semaphore([value])¶
یک شیء مشترک
threading.Semaphoreایجاد میکند و یک پراکسی برای آن برمیگرداند.
- Array(typecode, sequence)¶
یک آرایه ایجاد کنید و یک پراکسی برای آن برگردانید.
- Value(typecode, value)¶
یک شیء با ویژگی
valueقابل نوشتن ایجاد میکند و برای آن یک پراکسیبرمیگرداند.
- dict()¶
- dict(mapping)
- dict(sequence)
یک شیء
dictمشترک ایجاد میکند و یک پراکسی برای آن برمیگرداند.
- set()¶
- set(sequence)
- set(mapping)
یک شیء
setمشترک ایجاد کنید و یک پراکسی برای آن برگردانید.اضافه شده در نسخهی 3.14: پشتیبانی از
setافزوده شد.
تغییر یافته در نسخهی 3.6: اشیاء مشترک میتوانند تودرتو باشند. برای مثال، یک شیء ظرف مشترک مانند یک فهرست مشترک میتواند شامل اشیاء مشترک دیگر باشد که همگی توسط
SyncManagerمدیریت و همگامسازی میشوند.
- class multiprocessing.managers.Namespace¶
نوعی که میتواند در
SyncManagerثبت شود.یک شیء فضای نام هیچ متد عمومی ندارد، اما ویژگیهای قابلنوشتن دارد. نمایش آن مقادیر ویژگیهای آن را نشان میدهد.
با این حال، هنگام استفاده از یک پراکسی برای یک شیء فضای نام، ویژگیای که با
'_'شروع میشود، ویژگیای از پراکسی خواهد بود، نه ویژگیای از شیء مورد ارجاع:>>> mp_context = multiprocessing.get_context('spawn') >>> manager = mp_context.Manager() >>> Global = manager.Namespace() >>> Global.x = 10 >>> Global.y = 'hello' >>> Global._z = 12.3 # this is an attribute of the proxy >>> print(Global) Namespace(x=10, y='hello')
مدیرهای سفارشی¶
برای ایجاد مدیر اختصاصی خود، یک زیرکلاس از BaseManager ایجاد میکنید و از متد کلاسی register() برای ثبت انواع جدید یا callableها در کلاس مدیر استفاده میکنید. برای مثال:
from multiprocessing.managers import BaseManager
class MathsClass:
def add(self, x, y):
return x + y
def mul(self, x, y):
return x * y
class MyManager(BaseManager):
pass
MyManager.register('Maths', MathsClass)
if __name__ == '__main__':
with MyManager() as manager:
maths = manager.Maths()
print(maths.add(4, 3)) # prints 7
print(maths.mul(7, 8)) # prints 56
استفاده از مدیر راهدور¶
میتوان یک سرور مدیر را روی یک ماشین اجرا کرد و کاری کرد که کلاینتها بتوانند از ماشینهای دیگر از آن استفاده کنند (با فرض اینکه فایروالهای مرتبط اجازه این کار را بدهند).
اجرای دستورات زیر یک سرور برای یک صف مشترک واحد ایجاد میکند که کلاینتهای راه دور میتوانند به آن دسترسی داشته باشند:
>>> from multiprocessing.managers import BaseManager
>>> from queue import Queue
>>> queue = Queue()
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue', callable=lambda:queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()
یک کلاینت میتواند بهصورت زیر به سرور دسترسی یابد:
>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.put('hello')
یک کلاینت دیگر نیز میتواند از آن استفاده کند:
>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.get()
'hello'
فرایندهای محلی نیز میتوانند با استفاده از کد بالا روی کلاینت، به آن صف از راه دور دسترسی یابند:
>>> from multiprocessing import Process, Queue
>>> from multiprocessing.managers import BaseManager
>>> class Worker(Process):
... def __init__(self, q):
... self.q = q
... super().__init__()
... def run(self):
... self.q.put('local hello')
...
>>> queue = Queue()
>>> w = Worker(queue)
>>> w.start()
>>> class QueueManager(BaseManager): pass
...
>>> QueueManager.register('get_queue', callable=lambda: queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()
اشیای پراکسی¶
پراکسی شیءای است که به یک شیء اشتراکی ارجاع میدهد؛ این شیء اشتراکی (احتمالاً) در یک فرایند دیگر قرار دارد. به شیء اشتراکی، مورد ارجاع پراکسی گفته میشود. چندین شیء پراکسی میتوانند مورد ارجاع یکسانی داشته باشند.
یک شیء پراکسیدارای متدهایی است که متدهای متناظر مرجع خود را فراخوانی میکنند (اگرچه لزوماً هر متدی از مرجع از طریق پراکسی در دسترس نخواهد بود). به این ترتیب، میتوان از یک پراکسی دقیقاً مانند مرجع آن استفاده کرد:
>>> mp_context = multiprocessing.get_context('spawn')
>>> manager = mp_context.Manager()
>>> l = manager.list([i*i for i in range(10)])
>>> print(l)
[0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
>>> print(repr(l))
<ListProxy object, typeid 'list' at 0x...>
>>> l[4]
16
>>> l[2:5]
[4, 9, 16]
توجه داشته باشید که اعمال str() بر یک پراکسی بازنمایی مورد ارجاع را برمیگرداند، در حالی که اعمال repr() بازنمایی پراکسی را برمیگرداند.
یکی از ویژگیهای مهم اشیای پراکسی این است که آنها پیکلپذیر هستند، بنابراین میتوانند بین فرآیندها منتقل شوند. بههمین دلیل، یک مورد ارجاع میتواند شامل اشیای پراکسی باشد. این امر امکان تودرتو بودن این فهرستها، دیکشنریها و سایر اشیای پراکسی مدیریتشده را میدهد:
>>> a = manager.list()
>>> b = manager.list()
>>> a.append(b) # referent of a now contains referent of b
>>> print(a, b)
[<ListProxy object, typeid 'list' at ...>] []
>>> b.append('hello')
>>> print(a[0], b)
['hello'] ['hello']
بهطور مشابه، پراکسیهای دیکشنری و فهرست میتوانند درون یکدیگر تودرتو شوند:
>>> l_outer = manager.list([ manager.dict() for i in range(2) ])
>>> d_first_inner = l_outer[0]
>>> d_first_inner['a'] = 1
>>> d_first_inner['b'] = 2
>>> l_outer[1]['c'] = 3
>>> l_outer[1]['z'] = 26
>>> print(l_outer[0])
{'a': 1, 'b': 2}
>>> print(l_outer[1])
{'c': 3, 'z': 26}
اگر اشیاء استاندارد (غیرپراکسی) از نوع list یا dict در یک مورد ارجاع قرار داشته باشند، تغییرات آن مقادیر تغییرپذیر از طریق مدیر منتشر نمیشود، زیرا پراکسی هیچ راهی برای اطلاع از زمان تغییر مقادیر درون مورد ارجاع ندارد. با این حال، ذخیرهی یک مقدار در یک پراکسی ظرف (که باعث فراخوانی __setitem__ روی شیء پراکسی میشود) از طریق مدیر منتشر میشود؛ بنابراین برای تغییر مؤثر چنین آیتمی، میتوانید مقدار تغییرکرده را دوباره به پراکسی ظرف انتساب دهید:
# create a list proxy and append a mutable object (a dictionary)
lproxy = manager.list()
lproxy.append({})
# now mutate the dictionary
d = lproxy[0]
d['a'] = 1
d['b'] = 2
# at this point, the changes to d are not yet synced, but by
# updating the dictionary, the proxy is notified of the change
lproxy[0] = d
این رویکرد شاید برای بیشتر موارد استفاده، بهاندازهی بهکارگیری اشیای پراکسی تودرتو مناسب نباشد، اما همچنین سطحی از کنترل بر همگامسازی را نشان میدهد.
توجه
انواع پراکسی در multiprocessing هیچ کاری برای پشتیبانی از مقایسههای بر اساس مقدار انجام نمیدهند. بنابراین، برای مثال، داریم:
>>> manager.list([1,2,3]) == [1,2,3]
False
هنگام مقایسه، باید بهجای آن فقط از یک کپی از مورد ارجاع استفاده کرد.
- class multiprocessing.managers.BaseProxy¶
اشیای پراکسی نمونههایی از زیرکلاسهای
BaseProxyهستند.- _callmethod(methodname[, args[, kwds]])¶
متدی از مورد ارجاع پراکسی را فراخوانی کنید و نتیجهی آن را برگردانید.
اگر
proxyیک پراکسی باشد که مرجع آنobjاست، آنگاه عبارتproxy._callmethod(methodname, args, kwds)
عبارت را ارزیابی خواهد کرد
getattr(obj, methodname)(*args, **kwds)
در فرایند مدیر.
مقدار بازگشتی، یک کپی از نتیجه فراخوانی یا یک پراکسی به یک شیء مشترک جدید خواهد بود — مستندات آرگومان method_to_typeid در
BaseManager.register()را ببینید.اگر استثنایی از سوی فراخوانی پرتاب شود، سپس توسط
_callmethod()دوباره پرتاب میشود. اگر استثنای دیگری در فرایند مدیر پرتاب شود، آنگاه این استثنا به یک استثنایRemoteErrorتبدیل میشود و توسط_callmethod()پرتاب میشود.بهویژه توجه داشته باشید که اگر methodname در دسترس قرار نگرفته باشد، یک استثنا پرتاب خواهد شد.
مثالی از کاربرد
_callmethod():>>> l = manager.list(range(10)) >>> l._callmethod('__len__') 10 >>> l._callmethod('__getitem__', (slice(2, 7),)) # equivalent to l[2:7] [2, 3, 4, 5, 6] >>> l._callmethod('__getitem__', (20,)) # equivalent to l[20] Traceback (most recent call last): ... IndexError: list index out of range
- _getvalue()¶
یک کپی از مورد ارجاع برمیگرداند.
اگر مورد ارجاع غیرقابل پیکل (unpicklable) باشد، این عمل باعث پرتاب یک استثنا میشود.
- __repr__()¶
بازگرداندن بازنمایی از شیء پراکسی .
- __str__()¶
بازنمایی مورد ارجاع را برمیگرداند.
پاکسازی¶
یک شیء پراکسی از یک کالبک weakref استفاده میکند تا هنگامی که زبالهروبی میشود، خود را از ثبت مدیری که مالک مورد ارجاع آن است، خارج کند.
یک شیء مشترک زمانی از فرایند مدیر حذف میشود که دیگر هیچ پراکسیای به آن ارجاع نمیدهد.
استخرهای فرایند¶
میتوانید با استفاده از کلاس Pool استخری از فرایندها ایجاد کنید که وظایف ارسالشده به آن را انجام میدهد.
- class multiprocessing.pool.Pool([processes[, initializer[, initargs[, maxtasksperchild[, context]]]]])¶
یک شیء استخر فرایند (process pool) که استخری از فرایندهای کارگر را کنترل میکند و میتوان کارهایی را به آن ارسال کرد. این شیء از نتایج ناهمگام با مهلتهای زمانی و کالبکها پشتیبانی میکند و دارای پیادهسازی map موازی است.
processes تعداد فرایندهای کارگر مورد استفاده است. اگر processes برابر
Noneباشد، تعداد برگرداندهشده توسطos.process_cpu_count()استفاده میشود.اگر initializer برابر
Noneنباشد، هر فرایند کارگر در زمان شروع،initializer(*initargs)را فراخوانی خواهد کرد.maxtasksperchild تعداد وظایفی است که یک فرایند کارگر (worker process) میتواند پیش از آن که خارج شود و با یک فرایند کارگر تازه جایگزین شود به انجام برساند، تا امکان آزادسازی منابع استفادهنشده فراهم شود. مقدار پیشفرض maxtasksperchild برابر
Noneاست، به این معنا که فرایندهای کارگر تا زمانی که استخر وجود دارد زنده خواهند ماند.از context میتوان برای مشخص کردن زمینهای که برای شروع فرایندهای کارگر استفاده میشود، استفاده کرد. معمولاً یک استخر با استفاده از تابع
multiprocessing.Pool()یا متدPool()یک شیء زمینه ایجاد میشود. در هر دو مورد، context بهطور مناسب تنظیم میشود. اگرNoneباشد، فراخوانی این تابع روش شروع سراسری کنونی را بهصورت یک اثر جانبی تنظیم خواهد کرد، اگر قبلاً تنظیم نشده باشد. تابعget_context()را ببینید.توجه داشته باشید که متدهای شیء استخر را فقط باید فرایندی فراخوانی کند که آن را ایجاد کرده است.
هشدار
اشیای
multiprocessing.poolمنابع داخلی دارند که باید (مانند هر منبع دیگری) با استفاده از استخر بهعنوان یک مدیر زمینه یا با فراخوانی دستیclose()وterminate()بهدرستی مدیریت شوند. عدم انجام این کار میتواند منجر به آویزان ماندن فرایند در هنگام نهاییسازی شود.توجه داشته باشید که اتکا به زبالهرو برای نابود کردن استخر صحیح نیست، زیرا CPython تضمین نمیکند که نهاییساز استخر فراخوانی شود (برای اطلاعات بیشتر
object.__del__()را ببینید).تغییر یافته در نسخهی 3.2: پارامتر maxtasksperchild افزوده شد.
تغییر یافته در نسخهی 3.4: پارامتر context افزوده شد.
تغییر یافته در نسخهی 3.13: processes بهطور پیشفرض بهجای
os.cpu_count()ازos.process_cpu_count()استفاده میکند.توجه
فرایندهای کارگر درون یک
Poolمعمولاً در تمام مدت صف کارِ Pool زنده میمانند. یک الگوی رایج در سیستمهای دیگر (مانند Apache، mod_wsgi و غیره) برای آزادسازی منابعِ در اختیار کارگرها این است که به یک کارگر درون یک استخر اجازه داده شود تنها مقدار مشخصی از کار را پیش از خروج تکمیل کند، سپس پاکسازی شود و فرایند جدیدی برای جایگزینی با فرایند قدیمی ایجاد شود. آرگومان maxtasksperchild برایPoolاین قابلیت را در اختیار کاربر نهایی قرار میدهد.- apply(func[, args[, kwds]])¶
func را با آرگومانهای args و آرگومانهای کلیدواژهای kwds فراخوانی میکند. این فراخوانی تا آماده شدن نتیجه مسدود میشود. با توجه به این مسدودسازی،
apply_async()برای انجام کار بهصورت موازی مناسبتر است. علاوه بر این، func تنها در یکی از کارگرهای استخر اجرا میشود.
- apply_async(func[, args[, kwds[, callback[, error_callback]]]])¶
گونهای از متد
apply()که یک شیءAsyncResultرا برمیگرداند.اگر callback مشخص شده باشد، باید یک فراخوانیپذیر باشد که یک آرگومان میپذیرد. هنگامی که نتیجه آماده شود، callback روی آن اعمال میشود، یعنی مگر اینکه فراخوانی ناموفق باشد؛ در این صورت بهجای آن، error_callback اعمال میشود.
اگر error_callback مشخص شده باشد، باید یک شیء فراخوانیپذیر باشد که تنها یک آرگومان میپذیرد. اگر تابع هدف با شکست مواجه شود، error_callback با نمونهی استثنا فراخوانی میشود.
کالبکها باید بلافاصله به پایان برسند، در غیر این صورت نخی که نتایج را پردازش میکند مسدود خواهد شد.
- map(func, iterable[, chunksize])¶
معادل موازی تابع توکار
map()(هرچند فقط از یک آرگومان پیمایشپذیر پشتیبانی میکند؛ برای چند پیمایشپذیرstarmap()را ببینید). این متد تا آماده شدن نتیجه مسدود میشود.این متد، پیمایشپذیر را به تعدادی تکه تقسیم میکند و آنها را بهعنوان تکالیفجداگانه به استخر فرایند (process pool) ارسال میکند. اندازهی (تقریبی) این تکهها را میتوان با تنظیم chunksize روی یک عدد صحیح مثبت تعیین کرد.
توجه داشته باشید که این ممکن است برای پیمایشپذیرهای بسیار طولانی باعث مصرف بالای حافظه شود. برای کارایی بهتر، استفاده از
imap()یاimap_unordered()را با گزینهی صریح chunksize در نظر بگیرید.
- map_async(func, iterable[, chunksize[, callback[, error_callback]]])¶
گونهای از متد
map()که یک شیءAsyncResultرا برمیگرداند.اگر callback مشخص شده باشد، باید یک فراخوانیپذیر باشد که یک آرگومان میپذیرد. هنگامی که نتیجه آماده شود، callback روی آن اعمال میشود، یعنی مگر اینکه فراخوانی ناموفق باشد؛ در این صورت بهجای آن، error_callback اعمال میشود.
اگر error_callback مشخص شده باشد، باید یک شیء فراخوانیپذیر باشد که تنها یک آرگومان میپذیرد. اگر تابع هدف با شکست مواجه شود، error_callback با نمونهی استثنا فراخوانی میشود.
کالبکها باید بلافاصله به پایان برسند، در غیر این صورت نخی که نتایج را پردازش میکند مسدود خواهد شد.
- imap(func, iterable[, chunksize])¶
An iterator-based version of
map().آرگومان chunksize همان آرگومانی است که متد
map()از آن استفاده میکند. برای پیمایشپذیرهای بسیار طولانی، استفاده از مقدار بزرگ برای chunksize میتواند باعث تکمیل کار بسیار سریعتر از استفاده از مقدار پیشفرض1شود.همچنین اگر chunksize برابر
1باشد، متدnext()در پیمایشگر برگرداندهشده توسط متدimap()یک پارامتر اختیاری timeout دارد:next(timeout)در صورتی که نتیجه نتواند ظرف timeout ثانیه برگردانده شود،multiprocessing.TimeoutErrorرا پرتاب میکند.
- imap_unordered(func, iterable[, chunksize])¶
مشابه
imap()است، با این تفاوت که ترتیب نتایج حاصل از پیمایشگر برگرداندهشده باید دلخواه در نظر گرفته شود. (تنها زمانی که فقط یک فرایند کارگر وجود داشته باشد، تضمین میشود که ترتیب «صحیح» است.)
- starmap(func, iterable[, chunksize])¶
مانند
map()است، با این تفاوت که انتظار میرود عناصر iterable پیمایشپذیرهایی باشند که بهعنوان آرگومانها واگشایی میشوند.بنابراین، یک پیمایشپذیر از
[(1,2), (3, 4)]به[func(1,2), func(3,4)]منجر میشود.اضافه شده در نسخهی 3.3.
- starmap_async(func, iterable[, chunksize[, callback[, error_callback]]])¶
ترکیبی از
starmap()وmap_async()که یک iterable از پیمایشپذیرها را پیمایش میکند و func را با پیمایشپذیرهای واگشاییشده فراخوانی میکند. یک شیء نتیجه برمیگرداند.اضافه شده در نسخهی 3.3.
- close()¶
از ارسال وظایف بیشتر به استخر جلوگیری میکند. پس از تکمیل تمام وظایف، فرآیندهای کارگر خارج خواهند شد.
- terminate()¶
فرآیندهای کارگر را بلافاصله بدون تکمیل کارهای باقیمانده متوقف میکند. هنگامی که شیء استخر زبالهروبی شود،
terminate()بلافاصله فراخوانی خواهد شد.
- join()¶
منتظر بمانید تا فرآیندهای کارگر خارج شوند. پیش از استفاده از
join()بایدclose()یاterminate()را فراخوانی کنید.
تغییر یافته در نسخهی 3.3: اشیاء Pool اکنون از پروتکل مدیریت زمینه پشتیبانی میکنند -- Context Manager Types را ببینید.
__enter__()شیء Pool را بازمیگرداند، و__exit__()terminate()را فراخوانی میکند.
- class multiprocessing.pool.AsyncResult¶
کلاس نتیجهای که توسط
Pool.apply_async()وPool.map_async()بازگردانده میشود.- get([timeout])¶
نتیجه را به محض رسیدن برمیگرداند. اگر timeout
Noneنباشد و نتیجه ظرف timeout ثانیه نرسد،multiprocessing.TimeoutErrorپرتاب میشود. اگر فراخوانی راهدور استثنایی را پرتاب کرده باشد، آن استثنا توسطget()دوباره پرتاب میشود.
- wait([timeout])¶
تا زمانی که نتیجه در دسترس قرار بگیرد یا timeout ثانیه بگذرد، صبر کنید.
- ready()¶
برمیگرداند که آیا فراخوانی کامل شده است.
- successful()¶
برمیگرداند که آیا فراخوانی بدون پرتاب استثنا به پایان رسیده است یا خیر. اگر نتیجه آماده نباشد،
ValueErrorرا پرتاب میکند.تغییر یافته در نسخهی 3.7: اگر نتیجه آماده نباشد، به جای
AssertionError،ValueErrorپرتاب میشود.
مثال زیر استفاده از یک استخر را نشان میدهد:
from multiprocessing import Pool
import time
def f(x):
return x*x
if __name__ == '__main__':
with Pool(processes=4) as pool: # start 4 worker processes
result = pool.apply_async(f, (10,)) # evaluate "f(10)" asynchronously in a single process
print(result.get(timeout=1)) # prints "100" unless your computer is *very* slow
print(pool.map(f, range(10))) # prints "[0, 1, 4,..., 81]"
it = pool.imap(f, range(10))
print(next(it)) # prints "0"
print(next(it)) # prints "1"
print(it.next(timeout=1)) # prints "4" unless your computer is *very* slow
result = pool.apply_async(time.sleep, (10,))
print(result.get(timeout=1)) # raises multiprocessing.TimeoutError
شنوندهها و کلاینتها¶
معمولاً پیامرسانی بین فرایندها با استفاده از صفها یا اشیاء Connection برگرداندهشده توسط Pipe() انجام میشود.
با این حال، ماژول multiprocessing.connection انعطافپذیری بیشتری را فراهم میکند. این ماژول اساساً یک API پیاممحور سطح بالا برای کار با سوکتها یا پایپهای نامگذاریشده ویندوز ارائه میدهد. همچنین از احراز هویت digest با استفاده از ماژول hmac و پایش همزمان چندین اتصال پشتیبانی میکند.
- multiprocessing.connection.deliver_challenge(connection, authkey)¶
پیامی را که بهصورت تصادفی تولید شده است، به طرف دیگر اتصال بفرستید و منتظر پاسخ بمانید.
اگر پاسخ با خلاصهی پیام، با استفاده از authkey بهعنوان کلید، مطابقت داشته باشد، یک پیام خوشآمدگویی به سمت دیگر اتصال ارسال میشود. در غیر این صورت
AuthenticationErrorپرتاب میشود.
- multiprocessing.connection.answer_challenge(connection, authkey)¶
یک پیام را دریافت کنید، چکیده پیام را با استفاده از authkey بهعنوان کلید محاسبه کنید و سپس چکیده را بازفرستید.
اگر پیام خوشآمد دریافت نشود،
AuthenticationErrorپرتاب میشود.
- multiprocessing.connection.Client(address[, family[, authkey]])¶
تلاش میکند اتصالی به شنوندهای که از نشانی address استفاده میکند برقرار کند و یک
Connectionبرمیگرداند.نوع اتصال با آرگومان family تعیین میشود، اما عموماً میتوان این آرگومان را حذف کرد، زیرا معمولاً میتوان آن را از قالب address تشخیص داد. (به قالبهای نشانی مراجعه کنید)
اگر authkey داده شده باشد و
Noneنباشد، باید یک رشته بایت باشد و بهعنوان کلید مخفی برای یک چالش احراز هویت مبتنی بر HMAC استفاده خواهد شد. اگر authkey مقدارNoneباشد، هیچ احراز هویتی انجام نمیشود. در صورت شکست احراز هویت،AuthenticationErrorپرتاب میشود. کلیدهای احراز هویت را ببینید.
- class multiprocessing.connection.Listener([address[, family[, backlog[, authkey]]]])¶
پوششی برای یک سوکت مقیدشده یا پایپ نامدار ویندوز که برای اتصالها در حال «گوش دادن» است.
address نشانیای است که توسط سوکت مقیدشده یا پایپ نامدارِ شیء شنونده استفاده میشود.
توجه
اگر از نشانی '0.0.0.0' استفاده شود، این نشانی در ویندوز یک پایانه قابلاتصال نخواهد بود. اگر به یک پایانه قابلاتصال نیاز دارید، باید از '127.0.0.1' استفاده کنید.
family نوع سوکت (یا پایپ نامدار) مورد استفاده است. این مقدار میتواند یکی از رشتههای
'AF_INET'(برای سوکت TCP)،'AF_UNIX'(برای سوکت دامنهی یونیکس) یا'AF_PIPE'(برای پایپ نامدار ویندوز) باشد. از میان این موارد، تنها مورد اول تضمینشده در دسترس است. اگر family برابرNoneباشد، خانواده از قالب address استنتاج میشود. اگر address نیزNoneباشد، یک مقدار پیشفرض انتخاب میشود. این مقدار پیشفرض، خانوادهای است که فرض میشود سریعترین گزینهی در دسترس است. به قالبهای نشانی مراجعه کنید. توجه داشته باشید که اگر family برابر'AF_UNIX'و address برابرNoneباشد، سوکت در یک پوشهی موقت خصوصی ایجاد میشود که با استفاده ازtempfile.mkstemp()ایجاد شده است.اگر شیء شنونده از یک سوکت استفاده کند، backlog (بهطور پیشفرض ۱) پس از مقید شدن سوکت، به متد
listen()آن سوکت ارسال میشود.اگر authkey داده شده باشد و
Noneنباشد، باید یک رشته بایت باشد و بهعنوان کلید مخفی برای یک چالش احراز هویت مبتنی بر HMAC استفاده خواهد شد. اگر authkey مقدارNoneباشد، هیچ احراز هویتی انجام نمیشود. در صورت شکست احراز هویت،AuthenticationErrorپرتاب میشود. کلیدهای احراز هویت را ببینید.- accept()¶
یک اتصال را روی سوکت مقیدشده یا پایپ نامدارِ شیء شنونده بپذیرید و یک شیء
Connectionرا برگردانید. اگر احراز هویت انجام شود و ناموفق باشد، آنگاهAuthenticationErrorپرتاب میشود.
- close()¶
سوکت مقیدشده یا پایپ نامدار شیء شنونده را ببندید. این متد بهطور خودکار هنگامی که شنونده زبالهروبی میشود، فراخوانی میشود. با این حال، توصیه میشود آن را بهصورت صریح فراخوانی کنید.
اشیای Listener دارای ویژگیهای فقطخواندنی زیر هستند:
- address¶
نشانی که شیء Listener از آن استفاده میکند.
- last_accepted¶
نشانی که آخرین اتصال پذیرفتهشده از آن آمده است. اگر این نشانی در دسترس نباشد،
Noneخواهد بود.
تغییر یافته در نسخهی 3.3: اشیای شنونده اکنون از پروتکل مدیریت زمینه پشتیبانی میکنند — Context Manager Types را ببینید.
__enter__()شیء شنونده را برمیگرداند و__exit__()close()را فراخوانی میکند.
- multiprocessing.connection.wait(object_list, timeout=None)¶
تا آماده شدن یک شیء در object_list صبر میکند. فهرست اشیای موجود در object_list را که آمادهاند برمیگرداند. اگر timeout یک عدد اعشاری باشد، فراخوانی حداکثر به مدت آن تعداد ثانیه مسدود میشود. اگر timeout برابر
Noneباشد، برای مدت نامحدودی مسدود میشود. مهلت زمانی منفی معادل مهلت زمانی صفر است.در هر دو POSIX و Windows، یک شیء میتواند در object_list ظاهر شود اگر
یک شیء
Connectionقابلخواندن؛یک شیء
socket.socketمتصل و قابلخواندن؛ یا
یک شیء اتصال یا سوکت زمانی آماده است که دادهای برای خواندن از آن در دسترس باشد، یا طرف دیگر بسته شده باشد.
POSIX:
wait(object_list, timeout)تقریباً معادلselect.select(object_list, [], [], timeout)است. تفاوت این است که اگرselect.select()توسط یک سیگنال قطع شود، ممکن استOSErrorرا با کد خطایEINTRپرتاب کند، در حالی کهwait()این کار را نمیکند.ویندوز: یک آیتم در object_list باید یا یک دسته از نوع عدد صحیح باشد که انتظارپذیر است (بر اساس تعریف بهکاررفته در مستندات تابع Win32 یعنی
WaitForMultipleObjects()) یا میتواند شیء دارای متدfileno()باشد که یک دسته سوکت یا دسته پایپ را بازمیگرداند. (توجه داشته باشید که دستههای پایپ و دستههای سوکت، دستههای انتظارپذیر نیستند.)اضافه شده در نسخهی 3.3.
مثالها
کد سرور زیر یک شنونده ایجاد میکند که از 'secret password' بهعنوان کلید احراز هویت استفاده میکند. سپس منتظر یک اتصال میماند و مقداری داده به کلاینت ارسال میکند:
from multiprocessing.connection import Listener
from array import array
address = ('localhost', 6000) # family is deduced to be 'AF_INET'
with Listener(address, authkey=b'secret password') as listener:
with listener.accept() as conn:
print('connection accepted from', listener.last_accepted)
conn.send([2.25, None, 'junk', float])
conn.send_bytes(b'hello')
conn.send_bytes(array('i', [42, 1729]))
کد زیر به سرور متصل میشود و مقداری داده از سرور دریافت میکند:
from multiprocessing.connection import Client
from array import array
address = ('localhost', 6000)
with Client(address, authkey=b'secret password') as conn:
print(conn.recv()) # => [2.25, None, 'junk', float]
print(conn.recv_bytes()) # => 'hello'
arr = array('i', [0, 0, 0, 0, 0])
print(conn.recv_bytes_into(arr)) # => 8
print(arr) # => array('i', [42, 1729, 0, 0, 0])
کد زیر از wait() برای منتظر ماندن همزمان پیامها از چندین فرایند استفاده میکند:
from multiprocessing import Process, Pipe, current_process
from multiprocessing.connection import wait
def foo(w):
for i in range(10):
w.send((i, current_process().name))
w.close()
if __name__ == '__main__':
readers = []
for i in range(4):
r, w = Pipe(duplex=False)
readers.append(r)
p = Process(target=foo, args=(w,))
p.start()
# We close the writable end of the pipe now to be sure that
# p is the only process which owns a handle for it. This
# ensures that when p closes its handle for the writable end,
# wait() will promptly report the readable end as being ready.
w.close()
while readers:
for r in wait(readers):
try:
msg = r.recv()
except EOFError:
readers.remove(r)
else:
print(msg)
قالبهای نشانی¶
یک نشانی
'AF_INET'یک تاپلبه شکل(hostname, port)است که در آن hostname یک رشته و port یک عدد صحیح است.یک نشانی
'AF_UNIX'رشتهای است که نام یک پرونده در سامانه فایلبندی را نشان میدهد.نشانی
'AF_PIPE'رشتهای به شکلr'\\.\pipe\PipeName'است. برای استفاده ازClient()جهت اتصال به پایپای نامدار (named pipe) بر روی رایانهای راهدور به نام ServerName، باید بهجای آن از نشانیای به شکلr'\\ServerName\pipe\PipeName'استفاده کنید.
توجه داشته باشید که هر رشتهای که با دو بکاسلش آغاز شود، بهطور پیشفرض بهعنوان یک نشانی 'AF_PIPE' در نظر گرفته میشود، نه یک نشانی 'AF_UNIX'.
کلیدهای احراز هویت¶
هنگامی که از Connection.recv استفاده میکنید، دادههای دریافتی بهطور خودکار از پیکل خارج میشوند . متأسفانه پیکلگشایی دادهها از یک منبع نامطمئن، یک خطر امنیتی است. بنابراین Listener و Client() از ماژول hmac برای ارائه احراز هویت خلاصه (digest authentication) استفاده میکنند.
کلید احراز هویت یک رشتهی بایتی است که میتوان آن را بهعنوان یک گذرواژه در نظر گرفت: پس از برقراری اتصال، هر یک از دو طرف از طرف مقابل درخواست خواهد کرد که ثابت کند کلید احراز هویت را میداند. (اثبات اینکه هر دو طرف از یک کلید استفاده میکنند، شامل ارسال کلید از طریق اتصال نمیشود.)
اگر احراز هویت درخواست شده باشد اما هیچ کلید احراز هویتی مشخص نشده باشد، از مقدار بازگشتی current_process().authkey استفاده میشود (به Process مراجعه کنید). هر شیء Process که فرایند جاری ایجاد میکند، این مقدار را بهطور خودکار به ارث میبرد. این بدان معناست که (بهطور پیشفرض) همه فرایندهای یک برنامه چندفرایندی، یک کلید احراز هویت واحد را به اشتراک میگذارند که میتوان از آن هنگام برقراری اتصالها میان خودشان استفاده کرد.
همچنین میتوان کلیدهای احراز هویت مناسبی را با استفاده از os.urandom() تولید کرد.
این احراز هویت از اتصالهای Listener و Client() که از طریق نشانی قابل دسترسی هستند، محافظت میکند. این احراز هویت به پایپهای ناشناس ایجادشده توسط Pipe() یا پایپهایی که بهصورت داخلی توسط Queue استفاده میشوند، اعمال نمیشود. multiprocessing تمام فرآیندهای محلی را که با یک کاربر اجرا میشوند، مورد اعتماد میداند؛ در بیشتر سیستمعاملها، چنین فرآیندهایی در هر صورت میتوانند به توصیفگرهای پرونده پایپ یکدیگر دسترسی داشته باشند. برنامههایی که به جداسازی بین فرآیندهای یک کاربر نیاز دارند، باید آن را در سطح سیستمعامل فراهم کنند — برای مثال، با اجرای کارگرها تحت یک حساب کاربری متفاوت یا در یک سندباکس.
گزارشگیری¶
پشتیبانی تا حدی برای گزارشگیری در دسترس است. با این حال، توجه داشته باشید که بستهی logging از قفلهای مشترک بین فرایندها استفاده نمیکند، بنابراین ممکن است، بسته به نوع هندلر، پیامهای فرایندهای مختلف با هم مخلوط شوند.
- multiprocessing.get_logger()¶
گزارشگیر مورد استفادهی
multiprocessingرا برمیگرداند. در صورت نیاز، یک گزارشگیر جدید ایجاد خواهد شد.هنگامی که گزارشگیر برای اولین بار ایجاد میشود، سطح آن
logging.NOTSETاست و هیچ هندلر پیشفرضی ندارد. پیامهای ارسالشده به این گزارشگیر بهطور پیشفرض به گزارشگیر ریشه منتشر نمیشوند.توجه داشته باشید که در ویندوز، فرایندهای فرزند تنها سطح گزارشگیر فرایند والد را به ارث میبرند — هرگونه سفارشیسازی دیگری برای گزارشگیر به ارث برده نخواهد شد.
- multiprocessing.log_to_stderr(level=None)¶
این تابع،
get_logger()را فراخوانی میکند، اما علاوه بر بازگرداندن گزارشگیر ایجادشده توسط get_logger، یک هندلر اضافه میکند که خروجی را با استفاده از قالب'[%(levelname)s/%(processName)s] %(message)s'بهsys.stderrارسال میکند. شما میتوانیدlevelnameگزارشگیر را با ارسال یک آرگومانlevelتغییر دهید.
در ادامه یک نشست نمونه با فعال بودن گزارشدهی آمده است:
>>> import multiprocessing, logging
>>> logger = multiprocessing.log_to_stderr()
>>> logger.setLevel(logging.INFO)
>>> logger.warning('doomed')
[WARNING/MainProcess] doomed
>>> m = multiprocessing.Manager()
[INFO/SyncManager-...] child process calling self.run()
[INFO/SyncManager-...] created temp directory /.../pymp-...
[INFO/SyncManager-...] manager serving at '/.../listener-...'
>>> del m
[INFO/MainProcess] sending shutdown message to manager
[INFO/SyncManager-...] manager exiting with exitcode 0
برای مشاهدهی جدول کامل سطحهای گزارش، ماژول logging را ببینید.
ماژول multiprocessing.dummy¶
multiprocessing.dummy API ماژول multiprocessing را بازتولید میکند، اما چیزی بیش از دربرگیرندهای برای ماژول threading نیست.
بهویژه، تابع Pool ارائهشده توسط multiprocessing.dummy نمونهای از ThreadPool را برمیگرداند که زیرکلاسی از Pool است و از تمامی فراخوانیهای متد مشابه پشتیبانی میکند، اما بهجای فرایندهای کاری از استخری از نخهای کاری استفاده میکند.
- class multiprocessing.pool.ThreadPool([processes[, initializer[, initargs]]])¶
یک شیء استخر نخ که استخری از نخهای کارگر را کنترل میکند و میتوان کارها را به آن ارسال کرد. نمونههای
ThreadPoolاز نظر رابط کاملاً با نمونههایPoolسازگار هستند و منابع آنها نیز باید بهدرستی مدیریت شوند، چه با استفاده از استخر بهعنوان مدیر زمینه و چه با فراخوانی دستیclose()وterminate().processes تعداد نخهای کارگر مورد استفاده است. اگر processes برابر
Noneباشد، از تعداد برگرداندهشده توسطos.process_cpu_count()استفاده میشود.اگر initializer برابر
Noneنباشد، هر فرایند کارگر در زمان شروع،initializer(*initargs)را فراخوانی خواهد کرد.برخلاف
Pool، نمیتوان maxtasksperchild و context را ارائه کرد.توجه
ThreadPoolهمان رابطPoolرا دارد، که بر مبنای استخری از فرآیندها طراحی شده و پیش از معرفی ماژولconcurrent.futuresوجود داشته است. به همین دلیل، برخی عملیات را به ارث میبرد که برای استخری مبتنی بر نخها معنایی ندارند، و نوع خاص خود را برای نمایش وضعیت کارهای ناهمگام دارد،AsyncResult، که هیچ کتابخانه دیگری آن را درک نمیکند.کاربران بهطور کلی باید ترجیح دهند از
concurrent.futures.ThreadPoolExecutorاستفاده کنند، که رابط سادهتری دارد که از همان ابتدا بر پایهی نخها طراحی شده است و نمونههایconcurrent.futures.Futureرا برمیگرداند که با بسیاری از کتابخانههای دیگر، از جملهasyncioسازگار هستند.
دستورالعملهای برنامهنویسی¶
هنگام استفاده از multiprocessing، باید دستورالعملها و شیوههای رایج خاصی را رعایت کرد.
همهی روشهای شروع¶
موارد زیر برای تمام متدهای شروع صدق میکند.
از وضعیت مشترک اجتناب کنید
تا حد امکان، باید سعی کنید از جابهجایی حجم زیادی از دادهها بین فرآیندها خودداری کنید.
احتمالاً بهتر است برای ارتباط میان فرایندها، بهجای استفاده از سازوکارهای همگامسازی سطح پایینتر، از صفها یا پایپها استفاده کنید.
پیکلپذیری
اطمینان حاصل کنید که آرگومانهای متدهای پراکسیها پیکلپذیر باشند.
ایمنی نخی پراکسیها
از یک شیء پراکسی در بیش از یک نخ استفاده نکنید، مگر اینکه آن را با یک قفل محافظت کنید.
(هرگز مشکلی در استفادهی فرآیندهای مختلف از یک پراکسی یکسان وجود ندارد.)
پیوستن به فرآیندهای زامبی
در POSIX، هنگامی که یک فرایند پایان مییابد اما هنوز الحاق نشده باشد، به یک فرایند زامبی تبدیل میشود. هرگز نباید تعداد زیادی از آنها وجود داشته باشد، زیرا هر بار که یک فرایند جدید آغاز میشود (یا
active_children()فراخوانی میشود)، تمام فرایندهای پایانیافتهای که هنوز join نشدهاند، join خواهند شد. همچنین فراخوانیProcess.is_aliveبرای یک فرایند پایانیافته، آن فرایند را join میکند. با این حال، احتمالاً بهتر است تمام فرایندهایی را که آغاز میکنید، بهصراحت join کنید.
بهتر است به جای pickle/unpickle، ارثبری کنید
هنگام استفاده از روشهای آغاز spawn یا forkserver، بسیاری از انواع
multiprocessingباید پیکلپذیر باشند تا فرایندهای فرزند بتوانند از آنها استفاده کنند. با این حال، بهطور کلی باید از ارسال اشیاء مشترک به فرایندهای دیگر از طریق پایپها یا صفها خودداری کرد. در عوض، باید برنامه را بهگونهای تنظیم کنید که فرایندی که نیازمند دسترسی به یک منبع مشترک ایجادشده در جای دیگر است، بتواند آن را از یک فرایند نیاکان به ارث ببرد.
از خاتمه دادن به فرآیندها خودداری کنید
استفاده از متد
Process.terminateبرای توقف یک فرایند ممکن است باعث شود هر یک از منابع مشترک (مانند قفلها، سمافورها، پایپها و صفها) که در حال حاضر توسط آن فرایند استفاده میشوند، خراب یا برای سایر فرایندها غیرقابل دسترس شوند.بنابراین، احتمالاً بهتر است استفاده از
Process.terminateرا فقط برای فرایندهایی در نظر بگیرید که هرگز از هیچ منبع مشترکی استفاده نمیکنند.
پیوستن به فرآیندهایی که از صفها استفاده میکنند
به خاطر داشته باشید که فرایندی که آیتمهایی را در یک صف قرار داده است، پیش از خاتمه صبر میکند تا همه آیتمهای موجود در بافر توسط نخ «تغذیهکننده» به پایپ زیرین فرستاده شوند. (فرایند فرزند میتواند متد
Queue.cancel_join_threadصف را فراخوانی کند تا از این رفتار اجتناب کند.)این بدان معناست که هر زمان از یک صف استفاده میکنید، باید اطمینان حاصل کنید که همهی آیتمهایی که در صف قرار داده شدهاند، در نهایت پیش از آنکه فرآیند الحاق شود، از صف برداشته خواهند شد. در غیر این صورت، نمیتوانید مطمئن باشید که فرآیندهایی که آیتمهایی را در صف قرار دادهاند، خاتمه خواهند یافت. همچنین به یاد داشته باشید که فرآیندهای غیر daemon بهطور خودکار الحاق خواهند شد.
مثالی که دچار بنبست میشود بهصورت زیر است:
from multiprocessing import Process, Queue def f(q): q.put('X' * 1000000) if __name__ == '__main__': queue = Queue() p = Process(target=f, args=(queue,)) p.start() p.join() # this deadlocks obj = queue.get()یک راهحل در اینجا این است که دو خط آخر را جابهجا کنید (یا بهسادگی خط
p.join()را حذف کنید).
منابع را بهصراحت به فرآیندهای فرزند منتقل کنید
در POSIX، هنگام استفاده از روش راهاندازی fork، یک فرایند فرزند میتواند با استفاده از یک منبع سراسری، از یک منبع مشترک ایجادشده در فرایند والد استفاده کند. با این حال، بهتر است شیء را بهعنوان آرگومان به سازندهی فرایند فرزند ارسال کنید.
علاوه بر سازگار کردن (بهطور بالقوه) کد با ویندوز و سایر روشهای راهاندازی، این کار همچنین تضمین میکند که تا زمانی که فرایند فرزند هنوز زنده است، شیء در فرایند والد زبالهروبی نخواهد شد. این موضوع ممکن است در صورتی مهم باشد که با زبالهروبی شیء در فرایند والد، منبعی آزاد شود.
بنابراین برای مثال
from multiprocessing import Process, Lock def f(): ... do something using "lock" ... if __name__ == '__main__': lock = Lock() for i in range(10): Process(target=f).start()باید بهصورت زیر بازنویسی شود
from multiprocessing import Process, Lock def f(l): ... do something using "l" ... if __name__ == '__main__': lock = Lock() for i in range(10): Process(target=f, args=(lock,)).start()
مراقب جایگزینی sys.stdin با یک «شیء شبهپرونده» باشید
multiprocessingدر ابتدا بهصورت غیرشرطی فراخوانی میکرد:os.close(sys.stdin.fileno())در متد
multiprocessing.Process._bootstrap()--- این موضوع منجر به مشکلاتی در فرایندهای درونفرایند شد. این مورد به صورت زیر تغییر یافته است:sys.stdin.close() sys.stdin = open(os.open(os.devnull, os.O_RDONLY), closefd=False)این کار مشکل اساسی برخورد فرآیندها با یکدیگر را که به خطای توصیفگر پرونده نامعتبر منجر میشود، حل میکند، اما خطر بالقوهای برای برنامههایی ایجاد میکند که
sys.stdin()را با یک «شیء شبهپرونده» دارای بافرینگ خروجی جایگزین میکنند. این خطر آن است که اگر چندین فرآیندclose()را روی این شیءمانند پرونده فراخوانی کنند، ممکن است دادههای یکسان چندین بار به شیء تخلیه شوند و در نتیجه خرابی به وجود آید.اگر یک شیء شبهپرونده مینویسید و نهانگاهسازی خودتان را پیادهسازی میکنید، میتوانید با ذخیرهی pid هرگاه چیزی به نهانگاه اضافه میکنید و دور انداختن نهانگاه هنگامی که pid تغییر میکند، آن را ایمن در برابر انشعاب (fork-safe) کنید. برای مثال:
@property def cache(self): pid = os.getpid() if pid != self._pid: self._pid = pid self._cache = [] return self._cache
روشهای شروع spawn و forkserver¶
چند محدودیت اضافی وجود دارد که برای روش شروع fork اعمال نمیشود.
پیکلپذیری بیشتر (picklability)
اطمینان حاصل کنید که تمام آرگومانهای
Processپیکلپذیر هستند. همچنین، اگرProcess.__init__را زیرکلاسسازی میکنید، باید اطمینان حاصل کنید که نمونهها هنگام فراخوانی متدProcess.startپیکلپذیر خواهند بود.
متغیرهای سراسری
توجه داشته باشید که اگر کدی که در یک فرایند فرزند اجرا میشود، سعی کند به یک متغیر سراسری دسترسی پیدا کند، ممکن است مقداری که میبیند (در صورت وجود) با مقدار موجود در فرایند والد در زمانی که
Process.startفراخوانی شد، یکسان نباشد.با این حال، متغیرهای سراسری که صرفاً ثابتهای سطح ماژول هستند، مشکلی ایجاد نمیکنند.
ایمپورت ایمن ماژول اصلی
اطمینان حاصل کنید که ماژول اصلی میتواند بدون ایجاد عوارض جانبی ناخواسته (مانند آغاز یک فرایند جدید) بهصورت امن توسط یک مفسر پایتون جدید ایمپورت شود.
برای مثال، اجرای ماژول زیر با استفاده از متد شروع spawn یا forkserver با
RuntimeErrorشکست خواهد خورد:from multiprocessing import Process def foo(): print('hello') p = Process(target=foo) p.start()در عوض، باید از «نقطه ورود» برنامه با استفاده از
if __name__ == '__main__':به شکل زیر محافظت کرد:from multiprocessing import Process, freeze_support, set_start_method def foo(): print('hello') if __name__ == '__main__': freeze_support() set_start_method('spawn') p = Process(target=foo) p.start()(اگر برنامه بهجای اجرا بهصورت فریز بهصورت عادی اجرا شود، میتوان خط
freeze_support()را حذف کرد.)این به مفسر پایتون تازهایجادشده اجازه میدهد تا ماژول را بهصورت امن ایمپورت کند و سپس تابع
foo()آن ماژول را اجرا کند.محدودیتهای مشابهی در صورتی اعمال میشوند که یک استخر یا مدیر در ماژول اصلی ایجاد شود.
مثالها¶
نمایشی از چگونگی ایجاد و استفاده از مدیرها و پراکسیهای سفارشی:
from multiprocessing import freeze_support
from multiprocessing.managers import BaseManager, BaseProxy
import operator
##
class Foo:
def f(self):
print('you called Foo.f()')
def g(self):
print('you called Foo.g()')
def _h(self):
print('you called Foo._h()')
# A simple generator function
def baz():
for i in range(10):
yield i*i
# Proxy type for generator objects
class GeneratorProxy(BaseProxy):
_exposed_ = ['__next__']
def __iter__(self):
return self
def __next__(self):
return self._callmethod('__next__')
# Function to return the operator module
def get_operator_module():
return operator
##
class MyManager(BaseManager):
pass
# register the Foo class; make `f()` and `g()` accessible via proxy
MyManager.register('Foo1', Foo)
# register the Foo class; make `g()` and `_h()` accessible via proxy
MyManager.register('Foo2', Foo, exposed=('g', '_h'))
# register the generator function baz; use `GeneratorProxy` to make proxies
MyManager.register('baz', baz, proxytype=GeneratorProxy)
# register get_operator_module(); make public functions accessible via proxy
MyManager.register('operator', get_operator_module)
##
def test():
manager = MyManager()
manager.start()
print('-' * 20)
f1 = manager.Foo1()
f1.f()
f1.g()
assert not hasattr(f1, '_h')
assert sorted(f1._exposed_) == sorted(['f', 'g'])
print('-' * 20)
f2 = manager.Foo2()
f2.g()
f2._h()
assert not hasattr(f2, 'f')
assert sorted(f2._exposed_) == sorted(['g', '_h'])
print('-' * 20)
it = manager.baz()
for i in it:
print('<%d>' % i, end=' ')
print()
print('-' * 20)
op = manager.operator()
print('op.add(23, 45) =', op.add(23, 45))
print('op.pow(2, 94) =', op.pow(2, 94))
print('op._exposed_ =', op._exposed_)
##
if __name__ == '__main__':
freeze_support()
test()
استفاده از Pool:
import multiprocessing
import time
import random
import sys
#
# Functions used by test code
#
def calculate(func, args):
result = func(*args)
return '%s says that %s%s = %s' % (
multiprocessing.current_process().name,
func.__name__, args, result
)
def calculatestar(args):
return calculate(*args)
def mul(a, b):
time.sleep(0.5 * random.random())
return a * b
def plus(a, b):
time.sleep(0.5 * random.random())
return a + b
def f(x):
return 1.0 / (x - 5.0)
def pow3(x):
return x ** 3
def noop(x):
pass
#
# Test code
#
def test():
PROCESSES = 4
print('Creating pool with %d processes\n' % PROCESSES)
with multiprocessing.Pool(PROCESSES) as pool:
#
# Tests
#
TASKS = [(mul, (i, 7)) for i in range(10)] + \
[(plus, (i, 8)) for i in range(10)]
results = [pool.apply_async(calculate, t) for t in TASKS]
imap_it = pool.imap(calculatestar, TASKS)
imap_unordered_it = pool.imap_unordered(calculatestar, TASKS)
print('Ordered results using pool.apply_async():')
for r in results:
print('\t', r.get())
print()
print('Ordered results using pool.imap():')
for x in imap_it:
print('\t', x)
print()
print('Unordered results using pool.imap_unordered():')
for x in imap_unordered_it:
print('\t', x)
print()
print('Ordered results using pool.map() --- will block till complete:')
for x in pool.map(calculatestar, TASKS):
print('\t', x)
print()
#
# Test error handling
#
print('Testing error handling:')
try:
print(pool.apply(f, (5,)))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from pool.apply()')
else:
raise AssertionError('expected ZeroDivisionError')
try:
print(pool.map(f, list(range(10))))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from pool.map()')
else:
raise AssertionError('expected ZeroDivisionError')
try:
print(list(pool.imap(f, list(range(10)))))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from list(pool.imap())')
else:
raise AssertionError('expected ZeroDivisionError')
it = pool.imap(f, list(range(10)))
for i in range(10):
try:
x = next(it)
except ZeroDivisionError:
if i == 5:
pass
except StopIteration:
break
else:
if i == 5:
raise AssertionError('expected ZeroDivisionError')
assert i == 9
print('\tGot ZeroDivisionError as expected from IMapIterator.next()')
print()
#
# Testing timeouts
#
print('Testing ApplyResult.get() with timeout:', end=' ')
res = pool.apply_async(calculate, TASKS[0])
while 1:
sys.stdout.flush()
try:
sys.stdout.write('\n\t%s' % res.get(0.02))
break
except multiprocessing.TimeoutError:
sys.stdout.write('.')
print()
print()
print('Testing IMapIterator.next() with timeout:', end=' ')
it = pool.imap(calculatestar, TASKS)
while 1:
sys.stdout.flush()
try:
sys.stdout.write('\n\t%s' % it.next(0.02))
except StopIteration:
break
except multiprocessing.TimeoutError:
sys.stdout.write('.')
print()
print()
if __name__ == '__main__':
multiprocessing.freeze_support()
test()
مثالی برای نمایش نحوه استفاده از صفها جهت ارسال وظایف به مجموعهای از فرآیندهای کارگر و جمعآوری نتایج:
import time
import random
from multiprocessing import Process, Queue, current_process, freeze_support
#
# Function run by worker processes
#
def worker(input, output):
for func, args in iter(input.get, 'STOP'):
result = calculate(func, args)
output.put(result)
#
# Function used to calculate result
#
def calculate(func, args):
result = func(*args)
return '%s says that %s%s = %s' % \
(current_process().name, func.__name__, args, result)
#
# Functions referenced by tasks
#
def mul(a, b):
time.sleep(0.5*random.random())
return a * b
def plus(a, b):
time.sleep(0.5*random.random())
return a + b
#
#
#
def test():
NUMBER_OF_PROCESSES = 4
TASKS1 = [(mul, (i, 7)) for i in range(20)]
TASKS2 = [(plus, (i, 8)) for i in range(10)]
# Create queues
task_queue = Queue()
done_queue = Queue()
# Submit tasks
for task in TASKS1:
task_queue.put(task)
# Start worker processes
for i in range(NUMBER_OF_PROCESSES):
Process(target=worker, args=(task_queue, done_queue)).start()
# Get and print results
print('Unordered results:')
for i in range(len(TASKS1)):
print('\t', done_queue.get())
# Add more tasks using `put()`
for task in TASKS2:
task_queue.put(task)
# Get and print some more results
for i in range(len(TASKS2)):
print('\t', done_queue.get())
# Tell child processes to stop
for i in range(NUMBER_OF_PROCESSES):
task_queue.put('STOP')
if __name__ == '__main__':
freeze_support()
test()