Xây dựng REST API Bất Đồng Bộ Với FastAPI và Async SQLAlchemy 2.0: Tối Ưu Connection Pool & Migration Alembic

Development tutorial - IT technology blog
Development tutorial - IT technology blog

1. Bối cảnh & Vấn đề thực tế

Team mình từng gặp sự cố oái oăm trên production: API FastAPI viết chuẩn chỉnh nhưng khi chạm ngưỡng 1.500 req/s, latency bất ngờ vọt từ 25ms lên hơn 800ms. Nguyên nhân rất quen thuộc. Toàn bộ worker của Uvicorn bị nghẽn bởi driver database đồng bộ cũ (sync).

FastAPI chỉ phát huy tối đa hiệu năng khi toàn bộ pipeline I/O từ request đến database đều chạy non-blocking. SQLAlchemy 2.0 giải quyết triệt để bài toán này nhờ hỗ trợ async native qua driver asyncpg và chuẩn hóa cú pháp select().

Tuy nhiên, chuyển sang async không đơn giản là thêm từ khóa await. Bạn sẽ phải đối mặt với hai bẫy lớn:

  • Vô hiệu hóa Lazy-loading: Cơ chế ngầm nạp quan hệ (relationship) bị tắt hoàn toàn trong async. Gọi sai ngữ cảnh sẽ văng ngay lỗi MissingGreenlet.
  • Cạn kiệt Connection Pool: Cấu hình sai kích thước pool khiến kết nối bị treo ngầm, dẫn đến sập API khi traffic tăng đột biến.

2. Cài đặt & Cấu trúc thư mục

Khởi tạo virtual environment với Python 3.10+ và cài đặt các thư viện cần thiết:

# Khởi tạo môi trường ảo
python -m venv .venv
source .venv/bin/activate

# Cài đặt FastAPI, driver asyncpg và Alembic
pip install fastapi==0.110.0 uvicorn[standard]==0.29.0 \
    sqlalchemy==2.0.29 asyncpg==0.29.0 \
    alembic==1.13.1 pydantic-settings==2.2.1

Bố cục project theo hướng module hóa rõ ràng:

app/
├── api/
│   └── v1/
│       └── endpoints/
│           └── items.py
├── core/
│   ├── config.py
│   └── database.py
├── models/
│   └── item.py
├── schemas/
│   └── item.py
└── main.py
alembic/
├── env.py
└── versions/
alembic.ini

3. Thiết lập hệ thống

3.1. Cấu hình Async Engine & Connection Pool

Tại app/core/database.py, ta khởi tạo engine bất đồng bộ cùng các tham số pool quan trọng cho môi trường production:

from typing import AsyncGenerator
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession
from sqlalchemy.orm import DeclarativeBase

DATABASE_URL = "postgresql+asyncpg://postgres:postgres@localhost:5432/app_db"

engine = create_async_engine(
    DATABASE_URL,
    echo=False,
    pool_size=20,          # Giữ cố định 20 kết nối sẵn sàng trong pool
    max_overflow=10,       # Cho phép mở thêm tối đa 10 kết nối khi quá tải
    pool_timeout=30,       # Timeout 30s nếu không lấy được connection
    pool_recycle=1800,     # Tái tạo kết nối sau 30 phút, tránh ngắt kết nối ngầm
    pool_pre_ping=True,    # Ping kiểm tra kết nối sống trước khi cấp cho session
)

AsyncSessionLocal = async_sessionmaker(
    bind=engine,
    class_=AsyncSession,
    autoflush=False,
    expire_on_commit=False, # Tránh trigger lazy-load lỗi sau khi commit
)

class Base(DeclarativeBase):
    pass

async def get_db() -> AsyncGenerator[AsyncSession, None]:
    async with AsyncSessionLocal() as session:
        try:
            yield session
            await session.commit()
        except Exception:
            await session.rollback()
            raise

3.2. Khai báo Model chuẩn cú pháp 2.0

Trong app/models/item.py, sử dụng cú pháp Mapped và mapped_column để tận dụng khả năng suy luận kiểu dữ liệu (type hinting):

from sqlalchemy import String, Integer, DateTime, func
from sqlalchemy.orm import Mapped, mapped_column
from datetime import datetime
from app.core.database import Base

class Item(Base):
    __tablename__ = "items"

    id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True)
    title: Mapped[str] = mapped_column(String(255), nullable=False, index=True)
    description: Mapped[str | None] = mapped_column(String(500), nullable=True)
    created_at: Mapped[datetime] = mapped_column(
        DateTime(timezone=True), 
        server_default=func.now()
    )

3.3. Tích hợp Migration bất đồng bộ với Alembic

Khởi tạo thư mục cấu hình Alembic:

alembic init alembic

File alembic/env.py mặc định chỉ hỗ trợ synchronous engine. Cần cập nhật hàm run_migrations_online để sử dụng async engine với NullPool:

import asyncio
from logging.config import fileConfig
from sqlalchemy import pool
from sqlalchemy.engine import Connection
from sqlalchemy.ext.asyncio import async_engine_from_config
from alembic import context

from app.core.database import Base, DATABASE_URL
import app.models.item  # Đăng ký metadata của model vào Base

config = context.config
if config.config_file_name is not None:
    fileConfig(config.config_file_name)

target_metadata = Base.metadata
config.set_main_option("sqlalchemy.url", DATABASE_URL)

def do_run_migrations(connection: Connection) -> None:
    context.configure(connection=connection, target_metadata=target_metadata)
    with context.begin_transaction():
        context.run_migrations()

async def run_async_migrations() -> None:
    connectable = async_engine_from_config(
        config.get_section(config.config_ini_section, {}),
        prefix="sqlalchemy.",
        poolclass=pool.NullPool,
    )

    async with connectable.connect() as connection:
        await connection.run_sync(do_run_migrations)

    await connectable.dispose()

def run_migrations_online() -> None:
    asyncio.run(run_async_migrations())

if context.is_offline_mode():
    pass
else:
    run_migrations_online()

3.4. Xây dựng Endpoint API

Định nghĩa router trong app/api/v1/endpoints/items.py:

from fastapi import APIRouter, Depends, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select
from app.core.database import get_db
from app.models.item import Item
from pydantic import BaseModel

router = APIRouter(prefix="/items", tags=["items"])

class ItemCreate(BaseModel):
    title: str
    description: str | None = None

class ItemResponse(ItemCreate):
    id: int
    class Config:
        from_attributes = True

@router.post("/", response_model=ItemResponse, status_code=status.HTTP_201_CREATED)
async def create_item(payload: ItemCreate, db: AsyncSession = Depends(get_db)):
    new_item = Item(title=payload.title, description=payload.description)
    db.add(new_item)
    await db.flush()  # Gửi lệnh SQL xuống DB để lấy ID ngay mà không đóng transaction
    await db.refresh(new_item)
    return new_item

@router.get("/", response_model=list[ItemResponse])
async def list_items(skip: int = 0, limit: int = 20, db: AsyncSession = Depends(get_db)):
    query = select(Item).offset(skip).limit(limit)
    result = await db.execute(query)
    return result.scalars().all()

4. Kiểm thử & Vận hành thực chiến

Tạo và áp dụng bản migration đầu tiên:

alembic revision --autogenerate -m "init_items_table"
alembic upgrade head

Khởi chạy Uvicorn server:

uvicorn app.main:app --host 0.0.0.0 --port 8000 --reload

Để kiểm soát tình trạng rò rỉ kết nối trên production, hãy thêm endpoint giám sát trực tiếp Connection Pool:

@app.get("/health/db-pool")
async def pool_status():
    pool = engine.pool
    return {
        "pool_size": pool.size(),
        "checked_in": pool.checkedin(),
        "checked_out": pool.checkedout(),
        "overflow": pool.overflow(),
    }

Kinh nghiệm vận hành cần lưu ý:

  • Giải phóng session ngay lập tức: Tuyệt đối không gọi third-party API hoặc xử lý tác vụ CPU nặng bên trong block session. Hãy query xong dữ liệu rồi giải phóng kết nối trước khi xử lý tiếp.
  • Tính toán công thức Pool theo số lượng Worker: Uvicorn chạy 4 workers (-w 4) nghĩa là có 4 pool độc lập. Với pool_size=20 và max_overflow=10, hệ thống có thể chiếm tới 4 x 30 = 120 kết nối. Giá trị max_connections của PostgreSQL phải luôn lớn hơn con số này.
  • Bật pool_pre_ping=True: Tham số này tự động bắt và loại bỏ các stale connection bị ngắt bởi AWS RDS hoặc proxy mạng, hạn chế tối đa lỗi ConnectionDoesNotExistError.
Share: