آموزش ماژول asyncio – کلاس Queue

Please login to bookmark Close

این کد را ببینید؛ به‌نظرتان چه چاپ می‌کند؟ چند تولیدکننده داده می‌سازند و چند مصرف‌کننده باید آن را پردازش کنند، اما برنامه هرگز تمام نمی‌شود و در انتها معلق می‌ماند. چرا؟ چون هیچ ابزار امنی برای هماهنگی بین Taskها نداریم. اینجا دقیقاً جایی است که asyncio Queue وارد می‌شود: یک صف امن که بدون نیاز به قفل دستی، داده را بین Taskهای هم‌زمان جابه‌جا می‌کند.

import asyncio

async def consumer(name, items):
    for item in items:      # bug: shared list, no coordination
        print(name, item)

async def main():
    data = [1, 2, 3, 4]
    await asyncio.gather(
        consumer('c1', data),
        consumer('c2', data),  # both process the SAME items
    )

asyncio.run(main())

خروجی این کد نشان می‌دهد که هر دو مصرف‌کننده تمام آیتم‌ها را دوباره‌کاری می‌کنند؛ آیتم 2 دو بار پردازش می‌شود. هیچ‌کس نمی‌داند کدام آیتم قبلاً برداشته شده است. راه‌حل درست، گذاشتن داده در یک صف مشترک است تا هر آیتم فقط یک بار و توسط یک مصرف‌کننده برداشته شود.

چرا asyncio Queue و نه یک لیست ساده؟

یک list معمولی نه ظرفیت دارد، نه امکان انتظار. اگر مصرف‌کننده سریع‌تر از تولیدکننده باشد، باید بی‌کار بچرخد و CPU بسوزاند. asyncio.Queue این مشکل را حل می‌کند: متد get() وقتی صف خالی است Task را معلق می‌کند و به بقیهٔ کروتین‌ها اجازهٔ اجرا می‌دهد، و put() وقتی صف پر است منتظر خالی شدن جا می‌ماند. توجه کنید که رفتار آن شبیه queue.Queue در threading است، با این تفاوت که به‌جای مسدود کردن Thread، به‌صورت await کار می‌کند و با Event Loop سازگار است.

import asyncio

async def main():
    queue = asyncio.Queue(maxsize=5)  # 0 or unset = unbounded
    print(f'ظرفیت صف: {queue.maxsize}')

asyncio.run(main())

پارامتر maxsize اختیاری است؛ اگر صفر باشد یا تعیین نشود، صف نامحدود می‌شود. تعیین یک سقف در عمل مهم است: اگر تولیدکننده‌ها بسیار سریع‌تر از مصرف‌کننده‌ها باشند، صف نامحدود می‌تواند حافظه را ببلعد. با maxsize، تولیدکننده به‌طور طبیعی کند می‌شود تا مصرف‌کننده برسد؛ به این رفتار backpressure می‌گویند.

متدهای کلیدی

  • await queue.put(item): آیتم را اضافه می‌کند؛ اگر صف پر باشد تا خالی شدن جا منتظر می‌ماند.
  • queue.put_nowait(item): بدون انتظار اضافه می‌کند؛ اگر صف پر باشد QueueFull می‌دهد.
  • await queue.get(): یک آیتم برمی‌گرداند؛ اگر خالی باشد منتظر می‌ماند.
  • queue.task_done(): اعلام می‌کند پردازش یک آیتم گرفته‌شده تمام شده است.
  • await queue.join(): تا وقتی برای همهٔ آیتم‌ها task_done صدا زده شود منتظر می‌ماند.

مثال واقعی: الگوی Producer-Consumer

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

import asyncio
import random

async def producer(queue, name):
    for i in range(3):
        item = f'{name}-item-{i}'
        await asyncio.sleep(random.uniform(0.1, 0.5))
        await queue.put(item)          # waits if the queue is full
        print(f'{name}: تولید -> {item}')

async def consumer(queue, name):
    while True:                        # runs until cancelled
        item = await queue.get()       # suspends here when queue is empty
        await asyncio.sleep(random.uniform(0.1, 0.3))
        print(f'{name}: پردازش -> {item}')
        queue.task_done()              # signal one item is finished

async def main():
    queue = asyncio.Queue()
    producers = [asyncio.create_task(producer(queue, f'p-{i}')) for i in range(2)]
    consumers = [asyncio.create_task(consumer(queue, f'c-{i}')) for i in range(3)]

    await asyncio.gather(*producers)   # wait for all items to be produced
    await queue.join()                 # wait until every item is processed
    for c in consumers:
        c.cancel()                     # stop the infinite consumer loops

asyncio.run(main())

به‌یاد داشته باشید: چون مصرف‌کننده‌ها در حلقهٔ بی‌پایان‌اند، پس از queue.join() باید دستی با c.cancel() آن‌ها را متوقف کنید؛ در غیر این صورت برنامه هرگز پایان نمی‌یابد. این رایج‌ترین دامی است که تازه‌کارها در آن می‌افتند.

کی از این الگو استفاده کنیم؟

اگر فقط می‌خواهید تعداد Taskهای هم‌زمان را محدود کنید (مثلاً حداکثر ۱۰ درخواست هم‌زمان)، به‌جای Queue بهتر است سراغ Semaphore بروید که ساده‌تر است. اما وقتی جریان داده بین بخش‌های مختلف برنامه دارید و می‌خواهید تولید و مصرف را از هم جدا کنید، Queue انتخاب درست است. مزیت اصلی‌اش این است که تعداد تولیدکننده‌ها و مصرف‌کننده‌ها می‌توانند مستقل از هم تنظیم شوند.

جمع‌بندی

  • asyncio.Queue داده را امن و بدون قفل دستی بین Taskها منتقل می‌کند و متدهایش awaitپذیرند.
  • با maxsize یک سقف بگذارید تا backpressure داشته باشید و حافظه کنترل شود.
  • جفت task_done() و join() کلید پایان تمیز کار در الگوی Producer-Consumer است.
  • مصرف‌کننده‌های حلقه‌بی‌پایان را در انتها با cancel() متوقف کنید.
Please login to bookmark Close
پیشرفت شما در «دوره آموزش کانکارنسی در پایتون» (57%)
نظرات

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

57%
پیشرفت

سرفصل دوره

فهرست مطالب

سرفصل دوره

تمرین

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

پاسخ تمرین ها

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

اشتراک گذاری

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

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

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

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

تنظیمات

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