Tối ưu hiệu năng PostgreSQL trong Python với asyncpg: Xử lý nghẽn I/O và nạp 100k records/giây

Python tutorial - IT technology blog
Python tutorial - IT technology blog

1. Sự cố production: Backend sụp đổ khi truy vấn DB tăng đột biến

Khoảng 2h chiều thứ Ba tuần trước, hệ thống IoT bên mình bắt đầu đổ chuông cảnh báo PagerDuty liên tục. Latency API nhảy vọt từ 40ms lên 1.8s chỉ sau 5 phút. Gateway nhận dữ liệu từ 15.000 thiết bị cảm biến bị timeout hàng loạt, còn Redis queue thì dồn ứ hơn 200.000 message.

Phản xạ đầu tiên của team là kiểm tra database. Kỳ lạ thay, khi SSH vào cụm PostgreSQL (8 vCPU, 32GB RAM) chạy lệnh htop và check pg_stat_activity, CPU DB chỉ ăn 22%, RAM dư gần một nửa, disk write loanh quanh 15MB/s. Database hoàn toàn khỏe. Thủ phạm thực sự nằm ở tầng Python app: hàng trăm worker bị nghẽn mạng (I/O block) vì cạn connection pool mỗi khi query Postgres.

2. Vì sao psycopg2 truyền thống trở thành nút thắt cổ chai?

Dự án ban đầu dùng psycopg2 theo cơ chế đồng bộ (blocking I/O). Lưu lượng thấp tầm vài trăm req/s thì êm. Nhưng khi chạm ngưỡng 5.000 req/s, kiến trúc này gãy vì 3 điểm yếu cốt lõi:

  • Chặn thread xử lý (Thread-blocking): Mỗi lần gọi cursor.execute(), worker process phải đứng im đợi database phản hồi. Python không thể tận dụng khoảng nghỉ I/O này để nhận thêm việc khác.
  • Chi phí bắt tay kết nối (Connection Handshake): Mở mới một kết nối PostgreSQL tốn từ 30-50ms cho TLS handshake và xác thực user. Nếu app tạo kết nối liên tục, CPU server sẽ cạn kiệt chỉ để phục vụ kết nối rác.
  • Overhead từ Text Protocol: Nhiều driver cũ trao đổi dữ liệu dạng text thuần. Khi decode các kiểu dữ liệu nặng như JSONB, Timestamp hay UUID, CPU app phải parse chuỗi string liên tục, gây tụt throughput đáng kể.

3. Đặt 3 giải pháp lên bàn cân

Để cứu hệ thống, team mình đã thử nghiệm 3 phương án thực tế:

Phương án 1: Nâng số lượng Gunicorn / Uvicorn worker

Cách này dễ làm nhất nhưng tốn kém vô lý. Mỗi worker Python chiếm khoảng 110MB RAM. Tăng từ 8 lên 32 workers ngốn sạch gần 4GB RAM mà vẫn bị nghẽn mỗi khi DB có query chậm. Đây chỉ là giải pháp tình thế, không xử lý được gốc rễ vấn đề I/O.

Phương án 2: Dùng psycopg3 chế độ async

psycopg3 có hỗ trợ cú pháp async/await và giữ nguyên chuẩn DB-API 2.0. Nếu bạn đang refactor một codebase khổng lồ, đây là lựa chọn an toàn. Tuy nhiên, do phải duy trì tính tương thích ngược, tốc độ xử lý raw data của nó vẫn kém hơn các driver viết riêng cho AsyncIO.

Phương án 3: Chuyển hẳn sang asyncpg

asyncpg được MagicStack viết bằng Cython từ đầu, nhắm thẳng vào hiệu năng tối đa cho PostgreSQL. Thư viện này bỏ qua DB-API 2.0 để nói chuyện trực tiếp với Postgres qua binary protocol. Kết quả test thực tế: throughput tăng gấp 3.5 lần, còn CPU server app giảm gần 40% so với psycopg2.

4. Triển khai asyncpg: Từ kết nối cơ bản đến xử lý dữ liệu lớn

Cài đặt

Cài đặt thư viện qua pip:

pip install asyncpg

Khởi tạo Connection Pool dùng chung

Đừng bao giờ mở kết nối đơn lẻ trong từng endpoint. Hãy tạo một pool duy nhất khi app khởi động và dùng chung cho toàn bộ request:

import asyncio
import asyncpg

DATABASE_URL = "postgresql://postgres:[email protected]:5432/iot_db"

async def create_db_pool():
    return await asyncpg.create_pool(
        dsn=DATABASE_URL,
        min_size=10,        # Giữ sẵn 10 connection rảnh rỗi
        max_size=30,        # Tối đa 30 connection trên mỗi instance
        max_queries=50000,  # Tự động refresh connection sau 50k query để tránh leak bộ nhớ
        timeout=10.0        # Timeout nếu lấy connection từ pool quá 10 giây
    )

async def get_device_by_id(pool: asyncpg.Pool, device_id: str):
    # Mượn 1 connection từ pool, tự động trả lại khi thoát block async with
    async with pool.acquire() as conn:
        row = await conn.fetchrow(
            "SELECT id, name, firmware_version, created_at FROM devices WHERE id = $1",
            device_id
        )
        return dict(row) if row else None

Điểm cần nhớ: asyncpg dùng placeholder vị trí $1, $2, $3... thay vì %s hay ?.

Nạp dữ liệu lớn: Đừng dùng for loop, hãy dùng copy_records_to_table

Khi cần ghi 100.000 dòng dữ liệu cảm biến vào database, cách bạn viết code quyết định hệ thống sống hay chết. Chạy vòng lặp for rồi await conn.execute() từng dòng sẽ mất hơn 45 giây. Dùng executemany() giảm xuống khoảng 4.2 giây. Nhưng dùng copy_records_to_table (tận dụng lệnh COPY nhị phân của Postgres) chỉ mất 0.82 giây:

import time

async def bulk_insert_sensor_logs(pool: asyncpg.Pool, logs_data: list[tuple]):
    async with pool.acquire() as conn:
        start_time = time.perf_counter()
        
        # logs_data dạng: [('dev-001', 28.5, 65.2, 1718000000), ...]
        await conn.copy_records_to_table(
            table_name="sensor_logs",
            records=logs_data,
            columns=["device_id", "temperature", "humidity", "recorded_at"]
        )
        
        elapsed = time.perf_counter() - start_time
        print(f"Đã nạp {len(logs_data):,} bản ghi trong {elapsed:.2f}s")

Quản lý Transaction an toàn

Giao dịch tài chính hoặc trừ tồn kho cần đảm bảo tính toàn vẹn (ACID). Với asyncpg, bạn quản lý transaction cực kỳ gọn gàng qua context manager:

async def deduct_wallet_balance(pool: asyncpg.Pool, user_id: int, amount: float):
    async with pool.acquire() as conn:
        async with conn.transaction():
            # Khóa dòng để tránh race condition
            balance = await conn.fetchval(
                "SELECT balance FROM wallets WHERE user_id = $1 FOR UPDATE",
                user_id
            )
            if balance is None or balance < amount:
                raise ValueError("Số dư không đủ")
            
            await conn.execute(
                "UPDATE wallets SET balance = balance - $1 WHERE user_id = $2",
                amount, user_id
            )
            await conn.execute(
                "INSERT INTO audit_logs (user_id, amount, action) VALUES ($1, $2, 'WITHDRAW')",
                user_id, amount
            )

5. Kinh nghiệm vận hành trên Production

  • Công thức tính pool size: Đừng bao giờ đặt max_size = 100 bừa bãi. Công thức thực tế nên dùng: Pool Size = ((CPU Cores * 2) + Disk Count) / Số lượng app instance. Với cụm 4 worker app chạy trên server Postgres 8 cores NVMe, mỗi worker chỉ cần pool từ 5-10 connection.
  • JSONB không cần json.dumps(): asyncpg tự decode thẳng cột JSONB thành Python dict/list và ngược lại qua binary protocol. Bỏ các lệnh json.loads() thủ công để tiết kiệm chu kỳ CPU.
  • Graceful Shutdown: Luôn lắng nghe signal SIGTERM và gọi await pool.close() trong lifespan của FastAPI/Sanic. Việc này giúp đóng toàn bộ socket kết nối sạch sẽ, tránh để lại các process mồ côi (idle in transaction) làm treo database.
Share: