این کد را ببینید؛ بهنظرتان چه چاپ میکند؟ چند تولیدکننده داده میسازند و چند مصرفکننده باید آن را پردازش کنند، اما برنامه هرگز تمام نمیشود و در انتها معلق میماند. چرا؟ چون هیچ ابزار امنی برای هماهنگی بین 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()متوقف کنید.