Bỏ qua để đến nội dung

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 async khô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()create_task()
  • gatherTaskGroup (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ới asyncio.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.

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 xong
except StopIteration as e:
print(e.value) # xin chào An

Event 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.

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 queue

Khi 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 loop

Gọ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 asyncio
import 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.0s

create_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ả.

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 á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 ở 3

Khi 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)
import asyncio
import 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())
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:

  1. Luôn ném lại CancelledError sau khi dọn dẹp. Nuốt nó (except: pass) làm TaskGroup, timeout() và việc tắt chương trình hoạt động sai. Từ 3.8, CancelledError kế thừa BaseException chính để except Exception: không vô tình bắt nó.
  2. Dùng try/finally hoặc async with để giải phóng tài nguyên - chúng chạy cả khi bị huỷ.
  3. 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ế.
import asyncio
import 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ùng await 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)

Khi bắt buộc phải gọi code blocking:

import asyncio
import hashlib
import time
from 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())
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.

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())

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 asyncio
import 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 - 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 asyncio
from contextlib import asynccontextmanager
@asynccontextmanager
async 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.

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
  1. Viết fetch_all(urls, limit) dùng TaskGroupSemaphore, 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.
  2. 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.Queue chứa “token” được nạp lại định kỳ).
  3. 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.
  • Coroutine là phép tính có thể tạm dừng; gọi hàm async chỉ 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.
  • await là tuần tự; muốn đồng thời phải dùng Task - ưu tiên TaskGroup.
  • 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_thread hoặc process pool.
  • contextvarsthreading.local của thế giới async.

Bài tiếp theo: Chọn threading, multiprocessing hay asyncio?