asyncio chuyên sâu: event loop, Task, TaskGroup, huỷ và timeout
asyncio cho phép một thread duy nhất xử lý hàng nghìn kết nối mạng đồng thời - thứ mà thread hay process khó làm được vì tốn tài nguyên. Nhưng nó cũng là phần Python dễ dùng sai nhất: quên await, vô tình chặn event loop, task “biến mất” cùng exception của nó…
Bài này không chỉ dạy cú pháp async/await mà giải thích cơ chế bên dưới, để bạn tự suy luận được code async của mình sẽ chạy thế nào.
Trong bài này, bạn sẽ học:
- Coroutine thực chất là gì và vì sao gọi hàm
asynckhông chạy nó - Event loop lập lịch công việc như thế nào
- Sự khác biệt quan trọng giữa
await coro()vàcreate_task() gathervàTaskGroup(structured concurrency, 3.11+)- Timeout, cancellation và cách dọn dẹp đúng
- Giới hạn đồng thời với
Semaphore, pipeline vớiasyncio.Queue - Xử lý code blocking:
to_thread,run_in_executor contextvars, async generator và async context manager
Yêu cầu: Python 3.11 trở lên.
Mô hình: một người phục vụ, nhiều bàn
Phần tiêu đề “Mô hình: một người phục vụ, nhiều bàn”Hãy tưởng tượng một nhà hàng:
- Thread: mỗi bàn một người phục vụ riêng. Đứng chờ khách chọn món thì người đó đứng không. 1000 bàn cần 1000 người.
- asyncio: một người phục vụ. Nhận order bàn 1, đưa vào bếp, trong lúc chờ bếp thì sang bàn 2, bàn 3… Khi món bàn 1 xong, bếp gọi, người phục vụ quay lại mang món ra.
Điều kiện để mô hình này chạy: người phục vụ không bao giờ đứng chờ một việc gì. Nếu anh ta ở lại bàn 1 chờ khách suy nghĩ 10 phút, cả nhà hàng đứng im. Đây chính là quy tắc số một của asyncio: không được chặn event loop.
Đó là cooperative multitasking (đa nhiệm hợp tác): mỗi coroutine tự nguyện nhường quyền ở mỗi await, khác với thread bị hệ điều hành ngắt bất kỳ lúc nào.
Coroutine là gì?
Phần tiêu đề “Coroutine là gì?”import asyncio
async def greet(name): print(f"xin chào {name}") return len(name)
result = greet("An")print(result) # <coroutine object greet at 0x...> - CHƯA chạy gì cả!print(type(result)) # <class 'coroutine'>result.close() # tránh cảnh báo "coroutine was never awaited"Gọi một hàm async def không thực thi thân hàm. Nó chỉ tạo ra một coroutine object - một “phép tính có thể tạm dừng”, về bản chất giống generator (xem Generators và Iterators). Thậm chí bạn có thể tự “chạy tay” nó bằng send():
async def greet(name): return f"xin chào {name}"
coro = greet("An")try: coro.send(None) # chạy tới await đầu tiên hoặc tới khi xongexcept StopIteration as e: print(e.value) # xin chào AnEvent loop về cơ bản làm đúng việc này: gọi send(None) để chạy coroutine tới điểm await tiếp theo, ghi nhớ nó đang chờ gì, rồi chuyển sang coroutine khác.
Event loop hoạt động ra sao?
Phần tiêu đề “Event loop hoạt động ra sao?”Phiên bản đơn giản hoá của một event loop:
while còn việc: 1. Chạy mọi callback đã sẵn sàng (ready queue) - mỗi callback thường là "tiếp tục coroutine X tới await tiếp theo" 2. Hỏi hệ điều hành (select/epoll/kqueue): "socket nào có dữ liệu? chờ tối đa tới khi timer gần nhất hết hạn" 3. Với mỗi socket sẵn sàng / timer hết hạn: đưa callback tương ứng vào ready queueKhi coroutine await một thao tác mạng, nó đăng ký “đánh thức tôi khi socket này có dữ liệu” rồi trả quyền cho loop. Chính vì chỉ có một thread, mọi đoạn code giữa hai lần await chạy liền mạch, không bị chen ngang - nên trong asyncio bạn hiếm khi cần khoá cho các thao tác đơn giản.
Điểm vào tiêu chuẩn:
import asyncio
async def main(): await asyncio.sleep(1) print("xong")
asyncio.run(main()) # tạo loop, chạy main() tới khi xong, dọn dẹp, đóng loopGọi asyncio.run() một lần ở điểm vào chương trình. Không gọi nó bên trong code async.
Bẫy lớn nhất: await không có nghĩa là “chạy song song”
Phần tiêu đề “Bẫy lớn nhất: await không có nghĩa là “chạy song song””import asyncioimport time
async def fetch(name, delay): await asyncio.sleep(delay) # giả lập request mạng return f"{name} xong"
async def main(): start = time.perf_counter() a = await fetch("A", 1) b = await fetch("B", 1) print(a, b, f"{time.perf_counter() - start:.1f}s") # 2.0s !
asyncio.run(main())Mất 2 giây, không phải 1. await nghĩa là “chờ cái này xong rồi mới đi tiếp”. Hai lệnh await liên tiếp chạy tuần tự - chỉ là trong lúc chờ, loop có thể làm việc khác (nhưng ở đây không có việc khác nào).
Muốn chạy đồng thời, bạn phải tạo Task:
async def main(): start = time.perf_counter() task_a = asyncio.create_task(fetch("A", 1)) # lên lịch chạy NGAY task_b = asyncio.create_task(fetch("B", 1)) a = await task_a b = await task_b print(a, b, f"{time.perf_counter() - start:.1f}s") # 1.0screate_task() bọc coroutine thành Task và đưa vào event loop để chạy song song với code hiện tại. Bạn await task sau để lấy kết quả.
gather và TaskGroup
Phần tiêu đề “gather và TaskGroup”asyncio.gather
Phần tiêu đề “asyncio.gather”import asyncio
async def fetch(i): await asyncio.sleep(0.1) return i * 10
async def main(): results = await asyncio.gather(*(fetch(i) for i in range(5))) print(results) # [0, 10, 20, 30, 40] - đúng thứ tự đầu vào
asyncio.run(main())Vấn đề của gather khi có lỗi:
- Mặc định, task đầu tiên lỗi → exception được ném ra ngay, nhưng các task còn lại vẫn tiếp tục chạy ngầm không ai quản lý.
- Với
return_exceptions=True, exception được trả về như một giá trị trong list - bạn phải tự kiểm tra từng phần tử.
TaskGroup - structured concurrency (3.11+)
Phần tiêu đề “TaskGroup - structured concurrency (3.11+)”TaskGroup áp dụng nguyên tắc structured concurrency: mọi task sinh ra trong một khối phải kết thúc trước khi ra khỏi khối đó. Không có task “mồ côi”.
import asyncio
async def fetch(i): await asyncio.sleep(0.1 * i) if i == 3: raise ValueError(f"lỗi ở {i}") return i
async def main(): try: async with asyncio.TaskGroup() as tg: tasks = [tg.create_task(fetch(i)) for i in range(6)] print([t.result() for t in tasks]) except* ValueError as eg: # except* (3.11+) xử lý ExceptionGroup for exc in eg.exceptions: print("bắt được:", exc)
asyncio.run(main())# bắt được: lỗi ở 3Khi một task trong nhóm lỗi, TaskGroup tự động huỷ mọi task còn lại (task 4, 5), chờ chúng dọn dẹp xong, rồi ném ExceptionGroup chứa mọi lỗi. Đây là lựa chọn mặc định nên dùng cho code mới.
gather |
TaskGroup |
|
|---|---|---|
| Thứ tự kết quả | theo đầu vào | tự lấy từ từng task |
| Khi một task lỗi | các task khác chạy tiếp ngầm | huỷ các task còn lại |
| Nhiều lỗi cùng lúc | chỉ thấy lỗi đầu tiên | ExceptionGroup chứa tất cả |
| Thêm task động trong lúc chạy | không | có (tg.create_task bất kỳ lúc nào trong khối) |
Xử lý kết quả ngay khi có: as_completed
Phần tiêu đề “Xử lý kết quả ngay khi có: as_completed”import asyncioimport random
async def fetch(i): await asyncio.sleep(random.random()) return i
async def main(): for next_done in asyncio.as_completed([fetch(i) for i in range(5)]): print("xong:", await next_done) # thứ tự hoàn thành, không phải thứ tự đầu vào
asyncio.run(main())Timeout
Phần tiêu đề “Timeout”import asyncio
async def slow(): await asyncio.sleep(10)
async def main(): try: async with asyncio.timeout(1.5): # 3.11+ await slow() except TimeoutError: print("quá thời gian")
asyncio.run(main())asyncio.timeout() có thể bọc nhiều lệnh await - timeout áp dụng cho cả khối. asyncio.wait_for(coro, timeout) là cách cũ cho một coroutine đơn lẻ.
Luôn đặt timeout cho mọi thao tác mạng. Một server không phản hồi có thể làm coroutine của bạn chờ mãi mãi.
Cancellation: huỷ task và dọn dẹp đúng cách
Phần tiêu đề “Cancellation: huỷ task và dọn dẹp đúng cách”Huỷ một task bằng task.cancel(): lần tới task đó chạy, một CancelledError được ném vào vị trí await mà nó đang dừng.
import asyncio
async def worker(): try: while True: print("đang làm việc...") await asyncio.sleep(0.4) except asyncio.CancelledError: print("nhận lệnh huỷ, dọn dẹp...") await asyncio.sleep(0.1) # vẫn được await khi dọn dẹp raise # QUAN TRỌNG: ném lại! finally: print("đóng tài nguyên")
async def main(): task = asyncio.create_task(worker()) await asyncio.sleep(1) task.cancel() try: await task except asyncio.CancelledError: print("task đã bị huỷ:", task.cancelled())
asyncio.run(main())Quy tắc:
- Luôn ném lại
CancelledErrorsau khi dọn dẹp. Nuốt nó (except: pass) làmTaskGroup,timeout()và việc tắt chương trình hoạt động sai. Từ 3.8,CancelledErrorkế thừaBaseExceptionchính đểexcept Exception:không vô tình bắt nó. - Dùng
try/finallyhoặcasync withđể giải phóng tài nguyên - chúng chạy cả khi bị huỷ. - Cần bảo vệ một thao tác quan trọng không bị huỷ giữa chừng (ví dụ ghi transaction)? Dùng
await asyncio.shield(op())- nhưng hãy dùng hạn chế.
Quy tắc số một: không chặn event loop
Phần tiêu đề “Quy tắc số một: không chặn event loop”import asyncioimport time
async def tick(): for _ in range(3): print("tick", time.strftime("%X")) await asyncio.sleep(0.5)
async def bad_handler(): time.sleep(2) # ❌ chặn cả event loop 2 giây!
async def main(): await asyncio.gather(tick(), bad_handler())
asyncio.run(main())Trong 2 giây time.sleep chạy, tick không in được gì - mọi kết nối khác của server cũng đứng im. Những thứ chặn loop hay gặp:
time.sleep()→ dùngawait asyncio.sleep()requests.get()→ dùng thư viện async:httpx.AsyncClient,aiohttp- Driver database đồng bộ (
psycopg2,sqlite3) →asyncpg,psycopg(async),aiosqlite - Đọc/ghi file lớn, tính toán CPU nặng (parse JSON khổng lồ, xử lý ảnh)
Chuyển code blocking ra thread hoặc process
Phần tiêu đề “Chuyển code blocking ra thread hoặc process”Khi bắt buộc phải gọi code blocking:
import asyncioimport hashlibimport timefrom concurrent.futures import ProcessPoolExecutor
def blocking_io(): time.sleep(1) # thư viện đồng bộ return "dữ liệu"
def cpu_heavy(n): return sum(i * i for i in range(n))
async def main(): # I/O blocking -> chạy trong thread pool mặc định data = await asyncio.to_thread(blocking_io)
# CPU nặng -> chạy trong process pool loop = asyncio.get_running_loop() with ProcessPoolExecutor() as pool: total = await loop.run_in_executor(pool, cpu_heavy, 10_000_000)
print(data, total)
if __name__ == "__main__": asyncio.run(main())Phát hiện code chặn loop: debug mode
Phần tiêu đề “Phát hiện code chặn loop: debug mode”asyncio.run(main(), debug=True)# hoặc biến môi trường PYTHONASYNCIODEBUG=1Ở debug mode, asyncio cảnh báo mỗi khi một callback chạy quá 100ms (Executing <Task ...> took 2.003 seconds) và cảnh báo coroutine bị quên await.
Giới hạn mức đồng thời
Phần tiêu đề “Giới hạn mức đồng thời”Tạo 10.000 task gọi API cùng lúc sẽ làm server bên kia (hoặc chính bạn) quá tải. Dùng Semaphore:
import asyncio
async def fetch(sem, i): async with sem: # tối đa 10 request cùng lúc await asyncio.sleep(0.1) # giả lập request return i
async def main(): sem = asyncio.Semaphore(10) async with asyncio.TaskGroup() as tg: tasks = [tg.create_task(fetch(sem, i)) for i in range(100)] print(len([t.result() for t in tasks])) # 100, mất ~1 giây (10 đợt × 0.1s)
asyncio.run(main())Pipeline với asyncio.Queue
Phần tiêu đề “Pipeline với asyncio.Queue”Với lượng công việc rất lớn hoặc không biết trước, mẫu worker pool + queue tốt hơn tạo một task cho mỗi việc (không tạo hàng triệu task cùng lúc, có back-pressure):
import asyncio
async def producer(queue): for url in (f"https://site/{i}" for i in range(20)): await queue.put(url) # chờ nếu queue đầy -> back-pressure
async def worker(name, queue, results): while True: url = await queue.get() try: await asyncio.sleep(0.05) # giả lập tải results.append((name, url)) finally: queue.task_done()
async def main(): queue = asyncio.Queue(maxsize=5) results = [] async with asyncio.TaskGroup() as tg: workers = [tg.create_task(worker(f"w{i}", queue, results)) for i in range(4)] await producer(queue) await queue.join() # chờ mọi item được xử lý for w in workers: w.cancel() # dừng các worker đang chờ get() print(len(results)) # 20
asyncio.run(main())contextvars: “biến cục bộ” cho từng request
Phần tiêu đề “contextvars: “biến cục bộ” cho từng request”Trong server async, hàng nghìn request chạy xen kẽ trên cùng một thread, nên threading.local không dùng được. contextvars giải quyết điều này - mỗi Task có một bản copy context riêng:
import asyncioimport contextvars
request_id = contextvars.ContextVar("request_id", default="-")
def log(msg): print(f"[{request_id.get()}] {msg}") # không cần truyền request_id qua mọi hàm
async def handle(rid): request_id.set(rid) log("bắt đầu") await asyncio.sleep(0.1) # các request khác chạy xen vào đây log("kết thúc") # vẫn đúng request_id của mình
async def main(): async with asyncio.TaskGroup() as tg: for rid in ["req-1", "req-2", "req-3"]: tg.create_task(handle(rid))
asyncio.run(main())Đây là cơ chế mà các framework như FastAPI/Starlette, thư viện logging, tracing (OpenTelemetry) dùng để gắn thông tin request vào log.
Async generator và async context manager
Phần tiêu đề “Async generator và async context manager”Async generator - luồng dữ liệu đến dần qua mạng:
import asyncio
async def ticker(n): for i in range(n): await asyncio.sleep(0.1) yield i
async def main(): async for value in ticker(3): print(value) squares = [v * v async for v in ticker(3)] # async comprehension print(squares)
asyncio.run(main())Async context manager - tài nguyên cần await khi mở/đóng (kết nối, session):
import asynciofrom contextlib import asynccontextmanager
@asynccontextmanagerasync def connect(host): print("mở kết nối", host) await asyncio.sleep(0.1) try: yield f"conn:{host}" finally: await asyncio.sleep(0.1) print("đóng kết nối", host)
async def main(): async with connect("db.local") as conn: print("dùng", conn)
asyncio.run(main())Xem thêm bài Context manager nâng cao.
Những lỗi thường gặp
Phần tiêu đề “Những lỗi thường gặp”| Lỗi | Triệu chứng | Cách sửa |
|---|---|---|
Quên await |
RuntimeWarning: coroutine ... was never awaited, không có gì xảy ra |
Thêm await; bật debug mode |
await liên tiếp |
Chạy tuần tự, không nhanh hơn | TaskGroup / create_task |
| Gọi code blocking | Cả server bị “đơ” | Thư viện async, to_thread |
| Không lưu Task | Task biến mất giữa chừng | Lưu tham chiếu / TaskGroup |
Nuốt CancelledError |
Không tắt được, timeout không hoạt động | Luôn raise lại |
| Không có timeout | Coroutine treo vĩnh viễn | asyncio.timeout() |
| Tạo hàng triệu task | Tốn RAM, quá tải server đích | Semaphore, worker + Queue |
Bài tập
Phần tiêu đề “Bài tập”- Viết
fetch_all(urls, limit)dùngTaskGroupvàSemaphore, mỗi request có timeout 3 giây; URL lỗi/timeout được ghi vào danh sách riêng thay vì làm hỏng cả nhóm. - Viết một “rate limiter” async cho phép tối đa N lời gọi mỗi giây (gợi ý:
asyncio.Queuechứa “token” được nạp lại định kỳ). - Viết server chat TCP đơn giản bằng
asyncio.start_server: mỗi client gửi dòng nào thì broadcast cho mọi client khác.
Kết luận
Phần tiêu đề “Kết luận”- Coroutine là phép tính có thể tạm dừng; gọi hàm
asyncchỉ tạo coroutine chứ chưa chạy. - Event loop chạy mọi thứ trên một thread; coroutine nhường quyền ở mỗi
await. awaitlà tuần tự; muốn đồng thời phải dùng Task - ưu tiênTaskGroup.- Mọi thao tác mạng cần timeout; khi bị huỷ, dọn dẹp rồi ném lại
CancelledError. - Không bao giờ chặn event loop; code blocking đưa sang
to_threadhoặc process pool. contextvarslàthreading.localcủa thế giới async.
Bài tiếp theo: Chọn threading, multiprocessing hay asyncio?