صف‌ها

کد منبع: Lib/asyncio/queues.py


صف‌های asyncio به‌گونه‌ای طراحی شده‌اند که شبیه کلاس‌های ماژول queue باشند. اگرچه صف‌های asyncio در برابر نخ ایمن نیستند، اما برای استفاده به‌طور خاص در کد ناهمگام/await طراحی شده‌اند.

توجه داشته باشید که متدهای صف‌های asyncio پارامتر timeout ندارند؛ برای انجام عملیات صف با مهلت زمانی، از تابع asyncio.wait_for() استفاده کنید.

همچنین بخش Examples در زیر را ببینید.

صف

class asyncio.Queue(maxsize=0)

یک صف اولین ورودی، اولین خروجی (FIFO).

اگر maxsize کوچک‌تر یا مساوی صفر باشد، اندازه صف بی‌نهایت است. اگر این مقدار یک عدد صحیح بزرگ‌تر از 0 باشد، await put() هنگامی که صف به maxsize برسد مسدود می‌شود تا زمانی که یک آیتم توسط get() حذف شود.

برخلاف ماژول queue مربوط به threading در کتابخانه استاندارد، اندازه صف همیشه مشخص است و می‌توان آن را با فراخوانی متد qsize() برگرداند.

تغییر یافته در نسخه‌ی 3.10: پارامتر loop حذف شد.

این کلاس نخ‌ایمن نیست.

maxsize

تعداد آیتم‌های مجاز در صف.

empty()

اگر صف خالی باشد، True را برمی‌گرداند، در غیر این صورت False.

full()

اگر maxsize آیتم در صف وجود داشته باشد، True برمی‌گرداند.

اگر صف با maxsize=0 (پیش‌فرض) مقداردهی اولیه شده باشد، آنگاه full() هرگز True را برنمی‌گرداند.

async get()

یک آیتم را از صف حذف کرده و برمی‌گرداند. اگر صف خالی باشد، تا زمانی که یک آیتم در دسترس باشد صبر می‌کند.

اگر صف خاموش شده باشد و خالی باشد، یا اگر صف بلافاصله خاموش شده باشد، استثنای QueueShutDown را پرتاب می‌کند.

get_nowait()

اگر یک آیتم بلافاصله در دسترس باشد، آن را برمی‌گرداند؛ در غیر این صورت QueueEmpty را پرتاب می‌کند.

اگر صف خاموش شده باشد و خالی باشد، استثنای QueueShutDown را پرتاب می‌کند.

async join()

مسدود می‌شود تا تمام آیتم‌های موجود در صف دریافت و پردازش شده باشند.

تعداد وظایف ناتمام هر بار که یک آیتم به صف اضافه شود، افزایش می‌یابد. این تعداد هر بار که یک هم‌روال مصرف‌کننده task_done() را فراخوانی کند تا نشان دهد آن آیتم دریافت شده و تمام کارهای مربوط به آن کامل شده است، کاهش می‌یابد. هنگامی که تعداد وظایف ناتمام به صفر برسد، join() از حالت مسدود خارج می‌شود.

async put(item)

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

اگر صف خاموش شده باشد، QueueShutDown را پرتاب می‌کند.

put_nowait(item)

یک آیتم را بدون مسدودسازی در صف قرار دهید.

اگر هیچ جای خالی‌ای بلافاصله در دسترس نباشد، QueueFull پرتاب می‌شود.

اگر صف خاموش شده باشد، QueueShutDown را پرتاب می‌کند.

qsize()

تعداد آیتم‌های موجود در صف را برمی‌گرداند.

shutdown(immediate=False)

نمونه‌ای از Queue را در حالت خاموشی قرار دهید.

صف دیگر نمی‌تواند بزرگ‌تر شود. فراخوانی‌های بعدی put() باعث پرتاب QueueShutDown می‌شوند. فراخوانندگانی که در حال حاضر برای put() مسدود شده‌اند، از حالت مسدود خارج خواهند شد و QueueShutDown در وظیفه‌ای که پیش‌تر در انتظار بود پرتاب خواهد شد.

اگر immediate نادرست باشد (پیش‌فرض)، صف می‌تواند به‌صورت عادی با فراخوانی‌های get() برای استخراج وظایفی که از قبل بارگذاری‌شده‌اند، تخلیه شود.

و اگر task_done() برای هر وظیفه‌ی باقی‌مانده فراخوانی شود، یک join() در انتظار به‌طور عادی رفع انسداد خواهد شد.

پس از خالی شدن صف، فراخوانی‌های بعدی به get()، QueueShutDown را پرتاب خواهند کرد.

اگر immediate برابر true باشد، صف بلافاصله خاتمه می‌یابد. صف تخلیه می‌شود تا کاملاً خالی شود و تعداد وظایف ناتمام به تعداد وظایف تخلیه‌شده کاهش می‌یابد. اگر تعداد وظایف ناتمام صفر باشد، فراخواننده‌های join() از حالت مسدود خارج می‌شوند. همچنین، فراخواننده‌های مسدودشده‌ی get() از حالت مسدود خارج می‌شوند و QueueShutDown را به دلیل خالی بودن صف پرتاب می‌کنند.

هنگام استفاده از join() با immediate تنظیم‌شده روی true احتیاط کنید. این کار حتی زمانی که هیچ کاری روی وظایف انجام نشده باشد، انتظار join را رفع می‌کند و ناوردای معمولِ پیوستن به یک صف را نقض می‌کند.

اضافه شده در نسخه‌ی 3.13.

task_done()

نشان می‌دهد که یک آیتم کاری که پیش‌تر در صف قرار گرفته است، کامل شده است.

توسط مصرف‌کنندگان صف استفاده می‌شود. به ازای هر get() که برای واکشی یک آیتم کاری استفاده شود، فراخوانی متعاقب task_done() به صف اطلاع می‌دهد که پردازش آن آیتم کاری کامل شده است.

اگر join() در حال حاضر مسدودکننده باشد، هنگامی که همه‌ی آیتم‌ها پردازش شده باشند، از سر گرفته می‌شود (به این معنا که برای هر آیتمی که با put() در صف قرار داده شده باشد، یک فراخوانی task_done() دریافت شده باشد).

اگر بیشتر از تعداد آیتم‌های قرار داده‌شده در صف فراخوانی شود، ValueError پرتاب می‌کند.

صف اولویت

class asyncio.PriorityQueue

گونه‌ای از Queue؛ آیتم‌ها را به ترتیب اولویت (ابتدا کمترین) بازیابی می‌کند.

ورودی‌ها معمولاً تاپل‌هایی به شکل (priority_number, data) هستند.

صف LIFO

class asyncio.LifoQueue

گونه‌ای از Queue که ابتدا جدیدترین ورودی‌های افزوده‌شده را بازیابی می‌کند (آخرین ورودی، اولین خروجی).

استثناها

exception asyncio.QueueEmpty

این استثنا زمانی پرتاب می‌شود که متد get_nowait() روی یک صف خالی فراخوانی شود.

exception asyncio.QueueFull

استثنایی که با فراخوانی متد put_nowait() روی صفی که به maxsize خود رسیده باشد، پرتاب می‌شود.

exception asyncio.QueueShutDown

استثنایی که هنگام فراخوانی put()، put_nowait()، get() یا get_nowait() روی صفی که خاموش شده است، پرتاب می‌شود.

اضافه شده در نسخه‌ی 3.13.

مثال‌ها

می‌توان از صف‌ها برای توزیع بار کاری بین چندین وظیفه هم‌رو استفاده کرد:

import asyncio
import random
import time


async def worker(name, queue):
    while True:
        # Get a "work item" out of the queue.
        sleep_for = await queue.get()

        # Sleep for the "sleep_for" seconds.
        await asyncio.sleep(sleep_for)

        # Notify the queue that the "work item" has been processed.
        queue.task_done()

        print(f'{name} has slept for {sleep_for:.2f} seconds')


async def main():
    # Create a queue that we will use to store our "workload".
    queue = asyncio.Queue()

    # Generate random timings and put them into the queue.
    total_sleep_time = 0
    for _ in range(20):
        sleep_for = random.uniform(0.05, 1.0)
        total_sleep_time += sleep_for
        queue.put_nowait(sleep_for)

    # Create three worker tasks to process the queue concurrently.
    tasks = []
    for i in range(3):
        task = asyncio.create_task(worker(f'worker-{i}', queue))
        tasks.append(task)

    # Wait until the queue is fully processed.
    started_at = time.monotonic()
    await queue.join()
    total_slept_for = time.monotonic() - started_at

    # Cancel our worker tasks.
    for task in tasks:
        task.cancel()
    # Wait until all worker tasks are cancelled.
    await asyncio.gather(*tasks, return_exceptions=True)

    print('====')
    print(f'3 workers slept in parallel for {total_slept_for:.2f} seconds')
    print(f'total expected sleep time: {total_sleep_time:.2f} seconds')


asyncio.run(main())