Python hiện đại 6 — Async & Concurrency cho I/O
Trong Tổng quan Python hiện đại cho Data, chúng ta đã điểm qua bức tranh công cụ. Bài này giải quyết một câu hỏi rất thực chiến trong pipeline dữ liệu: làm sao gọi hàng nghìn API, đọc/ghi nhiều file và truy vấn nhiều DB cùng lúc mà không phải chờ tuần tự? Câu trả lời không chỉ là "dùng async" — mà là chọn đúng mô hình song song cho đúng loại việc. Chọn sai, bạn sẽ thấy code async còn chậm hơn code tuần tự, hoặc thấy 32 core mà chỉ một core chạy.
GIL và mô hình song song trong Python
GIL (Global Interpreter Lock) là một khóa toàn cục trong CPython — trình thông dịch Python phổ biến nhất. Quy tắc cốt lõi: tại một thời điểm, chỉ một thread được thực thi bytecode Python. Dù bạn tạo 10 thread trên máy 10 core, phần code Python thuần vẫn chạy nối tiếp nhau, không thực sự song song trên nhiều core.
Nghe như một khuyết điểm chí mạng, nhưng cần hiểu đúng: GIL được nhả ra trong lúc thread chờ I/O (đọc mạng, đọc đĩa, chờ DB trả về) và trong nhiều phép toán của thư viện C (NumPy, nén dữ liệu...). Chính vì vậy, để chọn công cụ đúng ta phải phân loại công việc:
- I/O-bound — phần lớn thời gian là chờ: chờ API trả về, chờ DB, chờ đọc file. CPU rảnh rỗi trong lúc chờ. Đây là loại việc chiếm đa số trong pipeline dữ liệu (ingest, gọi API đối tác, đọc/ghi storage).
- CPU-bound — phần lớn thời gian là tính: parse JSON khổng lồ, hash, nén, tính toán số học nặng bằng Python thuần. CPU luôn bận.
Nguyên tắc chọn:
| Loại việc | GIL có cản không? | Công cụ đúng |
|---|---|---|
| I/O-bound, nhiều tác vụ chờ | Không (nhả khi chờ) | asyncio hoặc thread |
| CPU-bound, tính nặng bằng Python | Có (chặn) | process (multiprocessing / ProcessPoolExecutor) |
| CPU-bound nhưng nằm trong lib C (NumPy/Polars) | Thường không | thread cũng có thể tận dụng nhiều core |
Với I/O-bound, GIL gần như vô hại vì thread nào cũng dành thời gian chờ — nhả khóa cho thread khác chạy. Với CPU-bound thuần Python, thread vô dụng (vẫn nối tiếp vì GIL), phải dùng process: mỗi tiến trình có trình thông dịch và GIL riêng, chạy thật sự song song trên nhiều core, đổi lại tốn RAM hơn và phải serialize dữ liệu khi truyền qua lại.
Xu hướng mới: từ Python 3.13, CPython bắt đầu cung cấp bản free-threaded (no-GIL) dạng thử nghiệm (tùy chọn build riêng), mở đường cho thread song song thật cho cả CPU-bound. Đây vẫn là giai đoạn đầu, nhiều thư viện chưa tương thích và hiệu năng đơn luồng có thể giảm. Ở mức hiện tại, hãy coi đây là thông tin định hướng, còn quyết định kiến trúc vẫn dựa trên mô hình GIL truyền thống mô tả ở trên.
Asyncio: chạy đồng thời trên một thread
asyncio là cách tiếp cận đồng thời hợp tác (cooperative concurrency) trên một thread duy nhất. Thay vì nhiều thread bị GIL kìm, ta có một event loop điều phối nhiều coroutine: khi một coroutine gặp điểm chờ I/O, nó chủ động "nhường" quyền điều khiển cho loop, để loop chạy coroutine khác trong lúc chờ. Không có việc chạy song song thực sự trên CPU — chỉ là không ai ngồi không trong lúc chờ.
async def / await / event loop
async defđịnh nghĩa một coroutine — gọi nó không chạy ngay mà trả về một đối tượng coroutine.await exprlà điểm nhường: "chờ ở đây, trong lúc chờ hãy cho việc khác chạy". Chỉ đượcawaitbên trongasync def.asyncio.run(main())khởi tạo event loop, chạy coroutine gốc tới khi xong.
import asyncio
async def fetch_one(name: str) -> str:
await asyncio.sleep(1) # minh hoạ một tác vụ I/O mất 1 giây
return f"done {name}"
async def main():
result = await fetch_one("A")
print(result)
asyncio.run(main())
Đoạn trên vẫn tuần tự. Sức mạnh chỉ xuất hiện khi ta chạy nhiều coroutine đồng thời.
Chạy đồng thời: gather và TaskGroup
asyncio.gather nhận nhiều coroutine, lên lịch tất cả và chờ chúng cùng hoàn tất. Ba tác vụ mỗi cái 1 giây sẽ xong trong ~1 giây thay vì 3:
async def main():
results = await asyncio.gather(
fetch_one("A"), fetch_one("B"), fetch_one("C")
)
print(results) # ['done A', 'done B', 'done C']
Từ Python 3.11, asyncio.TaskGroup là cách hiện đại hơn, dùng cấu trúc async with. Ưu điểm lớn: nếu một task lỗi, cả nhóm được hủy sạch (structured concurrency) — không để lại task "mồ côi" chạy ngầm.
async def main():
async with asyncio.TaskGroup() as tg:
t1 = tg.create_task(fetch_one("A"))
t2 = tg.create_task(fetch_one("B"))
print(t1.result(), t2.result())
Giới hạn đồng thời bằng Semaphore
Chạy 5.000 request cùng lúc sẽ làm sập server đối tác hoặc bị chặn IP. asyncio.Semaphore(n) giới hạn tối đa n tác vụ chạy đồng thời — số còn lại xếp hàng chờ tới lượt:
sem = asyncio.Semaphore(20) # tối đa 20 request cùng lúc
async def fetch_limited(url):
async with sem:
return await fetch_one(url)
Timeout và hủy
Một API "treo" không được để kéo cả pipeline. asyncio.timeout (3.11+) đặt hạn chót; quá hạn sẽ ném TimeoutError và hủy tác vụ đang chờ:
async def main():
try:
async with asyncio.timeout(5):
await fetch_one("chậm")
except TimeoutError:
print("bỏ qua vì quá 5 giây")
Việc hủy (cancellation) là bản chất của asyncio: khi task bị hủy, một CancelledError được ném vào coroutine tại điểm await, cho phép dọn dẹp (đóng kết nối, ghi log) trước khi thoát.
Sơ đồ: event loop xử lý nhiều I/O đồng thời
Điểm mấu chốt: trong lúc cả ba task "nhường" để chờ mạng, event loop không ngồi không — nó luân phiên xử lý bất kỳ task nào có dữ liệu về. Một thread, nhiều việc chờ, không lãng phí thời gian rỗi.
Thư viện async cho data
asyncio chỉ hữu ích khi thư viện I/O bạn dùng hỗ trợ async. Dùng thư viện đồng bộ (sync) bên trong coroutine sẽ chặn event loop (xem phần cạm bẫy). Các lựa chọn phổ biến trong công việc dữ liệu:
- httpx — HTTP client hỗ trợ cả sync lẫn async, API gần giống
requests. DùngAsyncClientđể gọi nhiều endpoint song song. aiohttp là lựa chọn async lâu đời khác, hiệu năng cao cho khối lượng rất lớn. - asyncpg — driver PostgreSQL async hiệu năng cao; cho phép nhiều truy vấn chạy đồng thời và quản lý connection pool. (
psycopgphiên bản 3 cũng có chế độ async.) - aiofiles — đọc/ghi file không chặn event loop, hữu ích khi ghi hàng loạt file trong lúc vẫn gọi mạng.
Khi thư viện chỉ có bản sync (nhiều SDK, driver cũ), đừng gọi thẳng trong coroutine. Đẩy nó sang một thread bằng asyncio.to_thread (3.9+) — coroutine await kết quả trong khi hàm sync chạy ở thread riêng, không chặn loop:
import asyncio
def sync_heavy_call(x): # hàm chỉ có bản đồng bộ
...
async def main():
result = await asyncio.to_thread(sync_heavy_call, 42)
Với nhu cầu tinh chỉnh hơn (pool cố định, tái sử dụng), dùng loop.run_in_executor với một ThreadPoolExecutor (I/O) hoặc ProcessPoolExecutor (CPU nặng). Đây là cầu nối quan trọng để chèn code sync/CPU vào một chương trình asyncio mà không phá vỡ event loop.
concurrent.futures: thread pool và process pool
Nếu bạn không muốn viết lại toàn bộ theo async/await, concurrent.futures cho cách song song đơn giản hơn với hai executor:
ThreadPoolExecutor— nhóm thread, hợp cho I/O-bound. GIL được nhả khi chờ I/O nên nhiều thread cùng chờ mạng/đĩa là hiệu quả. Dễ áp dụng khi thư viện chỉ có bản sync.ProcessPoolExecutor— nhóm tiến trình, hợp cho CPU-bound. Mỗi process có GIL riêng, chạy thật song song trên nhiều core; đổi lại chi phí khởi tạo và serialize dữ liệu qua ranh giới tiến trình.
from concurrent.futures import ThreadPoolExecutor
def download(url):
... # gọi requests (sync)
with ThreadPoolExecutor(max_workers=20) as ex:
results = list(ex.map(download, urls))
Đổi ThreadPoolExecutor thành ProcessPoolExecutor là chuyển sang song song đa tiến trình cho việc tính nặng — nhưng chỉ đúng khi việc thực sự CPU-bound, và hàm phải "picklable" (serialize được).
So sánh async vs thread vs process
| Tiêu chí | asyncio | ThreadPool (thread) | ProcessPool (process) |
|---|---|---|---|
| Hợp với | I/O-bound số lượng rất lớn | I/O-bound, code sync sẵn có | CPU-bound thuần Python |
| Song song CPU thật | Không (1 thread) | Không (bị GIL) | Có (nhiều tiến trình) |
| Chi phí mỗi đơn vị | Rất nhẹ (coroutine) | Trung bình (thread) | Nặng (tiến trình + serialize) |
| Số tác vụ đồng thời | Hàng nghìn+ | Hàng chục–trăm | ~ số core |
| Chia sẻ dữ liệu | Chung bộ nhớ | Chung bộ nhớ | Phải copy/IPC |
| Độ phức tạp code | Cao (async toàn trình) | Thấp | Trung bình |
| Cần lib async? | Có | Không | Không |
Quy tắc nhanh: hàng nghìn cuộc gọi mạng → asyncio; vài chục việc I/O với code sync có sẵn → ThreadPool; tính toán nặng bằng Python → ProcessPool.
Cạm bẫy thường gặp
1. Chặn event loop. Sai lầm phổ biến nhất: gọi hàm sync chậm (như requests.get, time.sleep, pd.read_parquet file lớn) hoặc tính CPU nặng ngay trong coroutine. Vì event loop chỉ có một thread, mọi coroutine khác đứng hình trong lúc đó — toàn bộ lợi ích đồng thời biến mất. Luôn dùng lib async, hoặc đẩy sang asyncio.to_thread/executor.
2. Quá nhiều kết nối cùng lúc. gather 10.000 request một lượt sẽ cạn socket/RAM, làm sập đối tác hoặc bị rate-limit chặn. Luôn giới hạn bằng Semaphore, và tôn trọng rate limit của API (ví dụ giới hạn số request/giây).
3. Xử lý lỗi trong gather. Mặc định, gather gặp một lỗi sẽ ném ngay lỗi đó và các task khác vẫn chạy ngầm (dễ rò rỉ). Hai lối đi: gather(..., return_exceptions=True) để gom cả kết quả lẫn exception vào danh sách rồi tự xử; hoặc dùng TaskGroup để có hành vi hủy nhóm rõ ràng.
4. Backpressure. Khi tốc độ tạo việc nhanh hơn tốc độ xử lý (ví dụ đọc hàng triệu dòng rồi bắn request), hàng đợi phình ra và tràn RAM. Giải pháp: xử lý theo lô (batch), dùng Semaphore/queue có giới hạn kích thước để nguồn phải chờ khi hạ nguồn quá tải.
5. Trộn lẫn sync và async. Không gọi asyncio.run bên trong loop đang chạy; quên await khiến coroutine không bao giờ chạy.
Ứng dụng trong pipeline dữ liệu
Các bài toán data điển hình hưởng lợi trực tiếp từ concurrency I/O:
- Gọi hàng nghìn API song song — làm giàu dữ liệu (enrichment), tra cứu tỉ giá, gọi dịch vụ đối tác cho từng bản ghi.
- Crawl/thu thập nhiều trang, nhiều nguồn cùng lúc.
- I/O file/DB đồng thời trong bước ingest: đọc nhiều file, ghi nhiều partition, chèn dữ liệu vào nhiều bảng song song.
Xem thêm cách các mảnh này ghép vào một pipeline hoàn chỉnh trong Pipeline dữ liệu ngân hàng.
Ví dụ: fetch song song nhiều endpoint (minh hoạ)
Đoạn dưới gọi nhiều endpoint bằng httpx.AsyncClient, giới hạn đồng thời bằng Semaphore, đặt timeout, và gom lỗi thay vì để cả lô đổ vỡ. Đây là minh hoạ khung code, không phải để chạy trong sandbox:
import asyncio
import httpx
CONCURRENCY = 20
sem = asyncio.Semaphore(CONCURRENCY)
async def fetch(client: httpx.AsyncClient, url: str) -> dict:
async with sem: # giới hạn 20 request đồng thời
try:
async with asyncio.timeout(10): # timeout mỗi request
resp = await client.get(url)
resp.raise_for_status()
return {"url": url, "ok": True, "data": resp.json()}
except Exception as e: # gom lỗi để không đổ cả lô
return {"url": url, "ok": False, "error": str(e)}
async def fetch_all(urls: list[str]) -> list[dict]:
async with httpx.AsyncClient() as client:
tasks = [fetch(client, u) for u in urls]
return await asyncio.gather(*tasks) # chạy đồng thời, giữ thứ tự
if __name__ == "__main__":
urls = [f"https://api.example.com/rate/{i}" for i in range(5000)]
results = asyncio.run(fetch_all(urls))
ok = sum(1 for r in results if r["ok"])
print(f"Thành công {ok}/{len(results)}")
Mấu chốt kiến trúc: dùng chung một AsyncClient (tái sử dụng connection pool, tránh mở/đóng kết nối liên tục), Semaphore chặn quá tải, timeout cắt request treo, và try/except trong mỗi task để một endpoint hỏng không kéo sập cả 5.000 cuộc gọi.
Use case thực tế
Bối cảnh NCB. Một pipeline hằng ngày cần làm giàu ~8.000 bản ghi giao dịch ngoại tệ bằng tỉ giá và thông tin từ dịch vụ API của đối tác. Bản triển khai đầu tiên viết theo kiểu tuần tự bằng requests: gọi từng bản ghi một, mỗi cuộc gọi trung bình ~300 ms (chủ yếu là chờ mạng, CPU gần như rảnh — đây rõ ràng là I/O-bound).
- Trước (tuần tự): 8.000 × 0,3 s ≈ 2.400 giây ≈ 40 phút, một core bận chờ, phần lớn thời gian là "ngồi không".
Nhóm chuyển sang httpx.AsyncClient + asyncio.gather với Semaphore(20) (tôn trọng rate-limit của đối tác), timeout 10 giây mỗi request và gom lỗi như ví dụ trên.
- Sau (async, 20 đồng thời): với 20 request chạy song song, thời gian thực tế xuống còn quãng 2–3 phút cho cùng 8.000 bản ghi — nhanh hơn khoảng 15–20 lần. (Con số là ước lượng minh hoạ; kết quả thật phụ thuộc rate-limit và độ trễ của đối tác.)
Các quyết định quan trọng và vì sao:
- Không dùng ProcessPool. Việc này là I/O-bound; nhiều tiến trình chỉ tốn RAM mà không nhanh hơn — chờ mạng thì thêm core vô nghĩa.
- Semaphore = 20, không phải 8.000. Bắn hết một lượt sẽ bị đối tác chặn (HTTP 429) và có thể sập kết nối; giới hạn đồng thời là bắt buộc để bền vững.
- Gom lỗi thay vì
raise. Vài chục bản ghi lỗi rải rác không được phép hủy cả mẻ 8.000; các bản lỗi được ghi log để retry ở lần chạy sau. - Bước tính CPU tách riêng. Sau khi có dữ liệu, phần tổng hợp/tính toán nặng được để cho Polars/NumPy (chạy đa lõi trong lib C) hoặc
ProcessPoolExecutor— không nhét vào event loop.
Kết quả: một job trước đây choán gần một giờ nay gọn trong vài phút, giảm rủi ro trễ SLA.
Ghi nhớ
- Phân loại việc trước, chọn công cụ sau. I/O-bound → asyncio/thread; CPU-bound thuần Python → process. Chọn sai là nguồn gốc của "async mà vẫn chậm".
- GIL khiến thread không song song CPU-bound, nhưng vô hại với I/O vì được nhả khi chờ. Free-threaded (no-GIL) 3.13+ còn thử nghiệm — chưa dựa vào để ra quyết định kiến trúc.
- asyncio: một event loop, nhiều coroutine hợp tác nhường nhau tại
await.gather/TaskGroup(3.11+) chạy đồng thời;TaskGroupcho hủy nhóm sạch sẽ. - Luôn giới hạn đồng thời bằng
Semaphorevà tôn trọng rate-limit; đặt timeout; xử lý hủy và lỗi trong gather một cách chủ động. - Đừng chặn event loop bằng code sync/CPU nặng — dùng lib async (httpx, asyncpg, aiofiles) hoặc đẩy sang
asyncio.to_thread/executor. - concurrent.futures cho song song đơn giản không cần async:
ThreadPoolExecutor(I/O),ProcessPoolExecutor(CPU). - Coi chừng backpressure: xử lý theo lô, dùng hàng đợi giới hạn để nguồn không làm tràn hạ nguồn.
Nguồn tham khảo
- Python Documentation —
asyncio— Asynchronous I/O (bao gồm Coroutines and Tasks:gather,TaskGroup,Semaphore,timeout,to_thread) - Python Documentation —
concurrent.futures— Launching parallel tasks (ThreadPoolExecutor,ProcessPoolExecutor) - Python Documentation —
threading— Thread-based parallelism - Python Documentation —
multiprocessing— Process-based parallelism - PEP 703 — Making the Global Interpreter Lock Optional in CPython
- HTTPX Documentation — Async Support và asyncpg Documentation
- AnyIO Documentation và Trio Documentation (structured concurrency)
Bài viết liên quan
Vì sao Python là ngôn ngữ số một của data engineer: vai trò trong pipeline (ingest/transform/orchestrate), hệ sinh thái thư viện (pandas/polars/pyarrow/sqlalchemy), quản lý môi trường (venv/uv/poetry), và khi nào dùng Python vs SQL/Spark.
Biến script thành pipeline đáng tin cậy: cấu trúc project & packaging (uv/poetry), type hints & pydantic, kiểm thử với pytest, logging & cấu hình, đóng gói Docker, và tích hợp CI cho code dữ liệu.
Học cách tổ chức code Python: định nghĩa hàm với tham số vị trí/từ khoá/mặc định, *args/**kwargs, lambda và hàm bậc cao, closure, decorator, generator với yield. Đóng gói code thành module và package, cô lập thư viện bằng môi trường ảo venv, quản lý phụ thuộc với pip và requirements.txt để dự án tái lập được trên mọi máy.
Hướng dẫn OOP trong Python từ class/instance, kế thừa và super(), đa hình & duck typing, encapsulation tới dunder methods, @property, classmethod/staticmethod, dataclass và type hints (mypy). Kèm nguyên tắc clean code: đặt tên rõ nghĩa, hàm nhỏ, DRY, SOLID cùng chuẩn PEP8 với công cụ ruff/black.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!