برای مدیریت همزمان چند وظیفه با asyncio TaskGroup در پایتون، وظایف را داخل یک بلاک مدیریت زمینه بسازید تا همزمانی ساختارمند، کنترل خطا و لغو وظایف به شکل قابل پیش بینی انجام شود. TaskGroup ساختن و نگهداری Taskها را ساده می کند، اگر یک وظیفه خطا دهد، بقیه به طور ایمن لغو می شوند و پس از خروج از بلاک می توانید نتیجه هر Task را بخوانید.

TaskGroup دقیقا چه مساله ای را حل می کند؟

TaskGroup پیاده سازی همزمانی ساختارمند در asyncio است. یعنی:

  • عمر همه Taskهای فرزند به عمر بلاکی که در آن ساخته می شوند گره می خورد.
  • اگر یکی از Taskها استثنا ایجاد کند، بقیه Taskها لغو می شوند و استثنا به صورت گروهی گزارش می شود.
  • هیچ Task یتیم یا نشتی منابع باقی نمی ماند و مسیر کنترل شما روشن است.

این رویکرد در برابر ساخت Taskهای پراکنده با create_task یا تکیه صرف بر gather، امن تر و قابل پیش بینی تر است؛ مخصوصا وقتی با خطا، لغو یا محدودیت زمانی سر و کار دارید.

شروع سریع: اجرای همزمان و جمع آوری نتایج

در این مثال چند کار شبیه سازی شده I/O را موازی اجرا و نتایج را به ترتیب ساخت Taskها جمع آوری می کنیم.

import asyncio
from time import perf_counter

async def job(name: str, delay: float) -> str:
    await asyncio.sleep(delay)
    return f"{name} done after {delay}s"

async def main():
    started = perf_counter()
    async with asyncio.TaskGroup() as tg:
        tasks = [
            tg.create_task(job("A", 1.0)),
            tg.create_task(job("B", 0.5)),
            tg.create_task(job("C", 0.2)),
        ]
    # در این نقطه همه Taskها تمام یا لغو شده اند
    results = [t.result() for t in tasks]
    duration = perf_counter() - started
    print("results:", results)
    print(f"took ~{duration:.2f}s (should be close to max delay)")

asyncio.run(main())

چطور بررسی کنیم درست کار می کند؟ چون اجرای همزمان است، زمان کل باید نزدیک به بیشترین تاخیر (اینجا حدود 1 ثانیه) باشد، نه مجموع تاخیرها.

رفتار خطا و لغو در TaskGroup

پاسخ کوتاه: اگر یکی از Taskها استثنا بدهد، TaskGroup بقیه Taskها را لغو می کند و یک ExceptionGroup بالا می آید. شما می توانید آن را بگیرید و تصمیم بگیرید چه کنید.

import asyncio

async def may_fail(i: int) -> int:
    await asyncio.sleep(0.2 * i)
    if i == 3:
        raise RuntimeError("boom at 3")
    return i * 10

async def main():
    try:
        async with asyncio.TaskGroup() as tg:
            tasks = [tg.create_task(may_fail(i)) for i in range(5)]
    except* RuntimeError as eg:
        # eg یک ExceptionGroup است
        print("caught:", [e.args[0] for e in eg.exceptions])
        # در این حالت برخی Taskها لغو شده اند و نتیجه ندارند
        # به نتایج موفق فقط زمانی دسترسی دارید که Taskها قبل از خطا تمام شده باشند.
        for t in tasks:
            if t.cancelled():
                print("task cancelled")
            elif t.done() and not t.exception():
                print("ok:", t.result())

asyncio.run(main())

نکته های کلیدی:

  • برای اینکه خطای یک وظیفه کل گروه را از بین نبرد، خطا را داخل همان وظیفه مدیریت کنید و خروجی امن برگردانید.
  • یا گروه ها را لایه بندی کنید: یک TaskGroup برای کارهای حیاتی و یک TaskGroup جدا برای کارهای اختیاری که خطایشان نباید گروه اصلی را متوقف کند.
import asyncio

async def safe_job(i: int):
    try:
        await asyncio.sleep(0.1 * i)
        if i % 2 == 0:
            raise ValueError("non-critical")
        return {"ok": True, "value": i}
    except Exception as e:
        return {"ok": False, "error": repr(e)}

async def main():
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(safe_job(i)) for i in range(6)]
    results = [t.result() for t in tasks]  # هیچ استثنایی بالا نمی آید
    print(results)

asyncio.run(main())

اعمال محدودیت زمانی و لغو تمیز

برای جلوگیری از آویزان شدن عملیات، یک سقف زمانی اعمال کنید. اگر زمان تمام شود، Taskها لغو می شوند و TimeoutError بالا می آید.

گزینه 1: context زمان سنجی

import asyncio

async def slow(i):
    await asyncio.sleep(1 + 0.3 * i)
    return i

async def run_with_timeout():
    try:
        async with asyncio.timeout(1.2):  # زمان کل برای گروه
            async with asyncio.TaskGroup() as tg:
                tasks = [tg.create_task(slow(i)) for i in range(5)]
    except TimeoutError:
        print("timed out!")
        # در این نقطه همه Taskهای باقیمانده لغو شده اند

asyncio.run(run_with_timeout())

گزینه 2: پوشاندن گروه با wait_for

import asyncio

async def do_group():
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(asyncio.sleep(2)) for _ in range(3)]

async def main():
    try:
        await asyncio.wait_for(do_group(), timeout=1.0)
    except asyncio.TimeoutError:
        print("group timed out")

asyncio.run(main())

کنترل فشار: محدود کردن همروندی با Semaphore

اگر با تعداد زیادی ورودی سروکار دارید (مثل درخواست های HTTP)، همروندی را محدود کنید تا به سرویس یا سیستم فشار نیاورید.

import asyncio
import aiohttp

async def fetch(session: aiohttp.ClientSession, url: str) -> int:
    async with session.get(url) as resp:
        await resp.read()
        return resp.status

async def crawl(urls, limit: int = 10):
    sem = asyncio.Semaphore(limit)
    async with aiohttp.ClientSession() as session:
        async with asyncio.TaskGroup() as tg:
            tasks = []
            for url in urls:
                async def bounded(u=url):
                    async with sem:
                        return await fetch(session, u)
                tasks.append(tg.create_task(bounded()))
    return [t.result() for t in tasks]

# استفاده:
# statuses = asyncio.run(crawl(list_of_urls, limit=20))

سه نکته اجرایی:

  • ClientSession را بیرون از تابع fetch بسازید و بین درخواست ها به اشتراک بگذارید.
  • برای بستن تمیز اتصال ها از async with استفاده کنید.
  • اگر لازم است هر درخواست خطا را خودش مدیریت کند، الگوی safe_job را به کار ببرید.

TaskGroup در برابر gather و create_task

کدام را کی انتخاب کنیم؟ معیار اصلی شما ساختار، مدیریت خطا و خوانایی است.

گزینهویژگی شاخصرفتار خطامناسب برای
TaskGroupهمزمانی ساختارمند، عمر Taskها محدود به بلاکلغو بقیه Taskها و بالا آمدن ExceptionGroupکدهای تولیدی و حساس با نیاز به مدیریت تمیز منابع
asyncio.gatherساده برای جمع آوری نتایجبه طور پیش فرض با اولین خطا شکست می خورد؛ با return_exceptions=True خطاها به خروجی تبدیل می شونداسکریپت های کوتاه یا جایی که ساختار سفت لازم نیست
create_taskکنترل کاملا دستی روی عمر Taskبه عهده شماست که خطا و لغو را مدیریت کنیدموارد خاص، وظایف بلندمدت، یا ادغام با حلقه رخداد سفارشی

الگوهای مفید با TaskGroup

۱) نگاشت همزمان با حفظ ترتیب ورودی

import asyncio

async def amap(func, items):
    async with asyncio.TaskGroup() as tg:
        tasks = [tg.create_task(func(x)) for x in items]
    # حفظ ترتیب بر اساس چینش اولیه
    return [t.result() for t in tasks]

۲) گروه های تو در تو برای ایزوله سازی خطا

import asyncio

async def core_task():
    ...

async def optional_task():
    raise ValueError("non-critical")

async def pipeline():
    async with asyncio.TaskGroup() as tg:
        tg.create_task(core_task())
        try:
            async with asyncio.TaskGroup() as aux:
                aux.create_task(optional_task())
        except* ValueError:
            # اجازه نمی دهیم خطای اختیاری خط لوله اصلی را متوقف کند
            pass

۳) گذارش پیشرفت بدون بلاک کردن

import asyncio

async def producer(q: asyncio.Queue):
    for i in range(5):
        await asyncio.sleep(0.2)
        await q.put(i)
    await q.put(None)  # سیگنال پایان

async def consumer(q: asyncio.Queue):
    while True:
        item = await q.get()
        if item is None:
            break
        print("got", item)

async def main():
    q = asyncio.Queue()
    async with asyncio.TaskGroup() as tg:
        tg.create_task(producer(q))
        tg.create_task(consumer(q))

asyncio.run(main())

خطاهای رایج و راه های اجتناب

  • فراموش کردن گرفتن رفرنس Taskها: اگر بعدا به نتیجه نیاز دارید، خروجی tg.create_task را ذخیره کنید؛ در غیر این صورت بعد از خروج از بلاک نتیجه ای برای خواندن ندارید.
  • قاطی کردن بلاکینگ I/O با async: هر عملیات مسدود کننده را با asyncio.to_thread یا یک کتابخانه async جایگزین کنید؛ وگرنه کل روال متوقف می شود.
  • نادیده گرفتن لغو: در تابع های async، کد طولانی را قابل لغو نگه دارید (به عنوان مثال با await های دوره ای یا چک کردن asyncio.current_task().cancelled()).
  • catch همه استثناها در بالا: اگر همه استثناها را کورکورانه بگیرید، خطاهای برنامه نویسی پنهان می شوند. فقط استثناهای پیش بینی شده را مدیریت کنید.
  • عدم تعیین سقف همروندی: برای کار با منابع محدود (پایگاه داده، API خارجی) از Semaphore استفاده کنید.

چطور مطمئن شویم واقعا همزمان اجرا می شود؟

دو راه ساده:

  1. زمان کل را با perf_counter بسنجید و آن را با بیشترین زمان یک کار مقایسه کنید.
  2. مهر زمان در ابتدای هر کار چاپ کنید؛ اگر خروجی ها روی هم می افتند، اجرای همزمان دارید.
import asyncio
from time import perf_counter

async def step(name, delay):
    t0 = perf_counter()
    await asyncio.sleep(delay)
    print(name, "elapsed", f"{perf_counter()-t0:.2f}s")

async def main():
    async with asyncio.TaskGroup() as tg:
        for i, d in enumerate([0.6, 0.4, 0.2], 1):
            tg.create_task(step(f"T{i}", d))

asyncio.run(main())

گام بعدی چیست؟

TaskGroup زمانی می درخشد که اجرای موازی، مدیریت خطا و بستن تمیز منابع برایتان حیاتی است. قدم بعدی را با بازنگری بخش های دارای gather یا create_task در کدتان بردارید: اگر به ساختار، لغو قابل اتکا و عمر کنترل شده Taskها نیاز دارید، آن بخش ها را به TaskGroup مهاجرت دهید. سپس برای بارهای سنگین، محدودیت همروندی و زمان سنجی سراسری را اضافه کنید تا سامانه پایدار و قابل پیش بینی بماند.