Chuyển đến nội dung chính

第 16 課:文件上傳、快取和非同步深入探討

檔案上傳流、UploadFile、多部分錶單。使用 fastapi-cache2 進行 Redis 快取。非同步/等待深入研究:非同步、並發任務、信號量、非同步產生器。 httpx 用於非同步 HTTP 用戶端。

💻 程式設計 — 第 16 課 第 16 課:文件上傳、快取和深度非同步 潛水

Python FastAPI:從基礎到進階

第 4 部分:進階功能

亞洲開發網

1. 文件上傳

import os
import uuid
from pathlib import Path

from fastapi import APIRouter, File, HTTPException, UploadFile

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

UPLOAD_DIR = Path("uploads")
UPLOAD_DIR.mkdir(exist_ok=True)

ALLOWED_TYPES = {"image/jpeg", "image/png", "image/webp", "application/pdf"}
MAX_FILE_SIZE = 10 * 1024 * 1024  # 10MB


@router.post("/upload")
async def upload_file(file: UploadFile = File(...)):
    """Upload single file."""
    # Validate content type
    if file.content_type not in ALLOWED_TYPES:
        raise HTTPException(400, f"File type {file.content_type} not allowed")

    # Validate size
    contents = await file.read()
    if len(contents) > MAX_FILE_SIZE:
        raise HTTPException(400, "File too large (max 10MB)")

    # Generate unique filename
    ext = Path(file.filename or "file").suffix
    filename = f"{uuid.uuid4()}{ext}"
    filepath = UPLOAD_DIR / filename

    # Save file
    with open(filepath, "wb") as f:
        f.write(contents)

    return {
        "filename": filename,
        "original_name": file.filename,
        "content_type": file.content_type,
        "size": len(contents),
        "url": f"/uploads/{filename}",
    }


@router.post("/upload-streaming")
async def upload_large_file(file: UploadFile = File(...)):
    """Upload file lớn với streaming (không load toàn bộ vào memory)."""
    if file.content_type not in ALLOWED_TYPES:
        raise HTTPException(400, "File type not allowed")

    ext = Path(file.filename or "file").suffix
    filename = f"{uuid.uuid4()}{ext}"
    filepath = UPLOAD_DIR / filename

    total_size = 0
    chunk_size = 1024 * 1024  # 1MB chunks

    with open(filepath, "wb") as f:
        while chunk := await file.read(chunk_size):
            total_size += len(chunk)
            if total_size > MAX_FILE_SIZE:
                os.unlink(filepath)
                raise HTTPException(400, "File too large")
            f.write(chunk)

    return {"filename": filename, "size": total_size}


@router.post("/upload-multiple")
async def upload_multiple_files(files: list[UploadFile] = File(...)):
    """Upload nhiều files cùng lúc."""
    results = []
    for file in files[:10]:  # Max 10 files
        contents = await file.read()
        ext = Path(file.filename or "file").suffix
        filename = f"{uuid.uuid4()}{ext}"
        filepath = UPLOAD_DIR / filename

        with open(filepath, "wb") as f:
            f.write(contents)

        results.append({
            "filename": filename,
            "original_name": file.filename,
            "size": len(contents),
        })
    return results

2.Redis緩存

uv add fastapi-cache2[redis]
# app/core/cache.py
from fastapi_cache import FastAPICache
from fastapi_cache.backends.redis import RedisBackend
from fastapi_cache.decorator import cache
import redis.asyncio as redis

from app.config import settings


async def init_cache():
    """Khởi tạo Redis cache (gọi trong lifespan)."""
    redis_client = redis.from_url(
        settings.redis_url,
        encoding="utf-8",
        decode_responses=True,
    )
    FastAPICache.init(RedisBackend(redis_client), prefix="fastapi-cache:")


# Sử dụng trong routes
from fastapi import APIRouter

router = APIRouter()


@router.get("/items/{item_id}")
@cache(expire=60)  # Cache 60 giây
async def get_item(item_id: int):
    """Response được cache trong Redis."""
    # Expensive operation...
    return {"item_id": item_id, "data": "expensive_result"}


@router.get("/stats")
@cache(expire=300)  # Cache 5 phút
async def get_stats():
    """Statistics - cache lâu hơn."""
    return {"users": 1000, "posts": 5000}

自訂緩存助手

# app/core/cache_helper.py
import json
from typing import Any

import redis.asyncio as redis


class CacheService:
    def __init__(self, redis_client: redis.Redis):
        self.redis = redis_client

    async def get(self, key: str) -> Any | None:
        value = await self.redis.get(key)
        if value:
            return json.loads(value)
        return None

    async def set(self, key: str, value: Any, expire: int = 3600) -> None:
        await self.redis.setex(key, expire, json.dumps(value, default=str))

    async def delete(self, key: str) -> None:
        await self.redis.delete(key)

    async def delete_pattern(self, pattern: str) -> None:
        """Xóa tất cả keys theo pattern."""
        async for key in self.redis.scan_iter(match=pattern):
            await self.redis.delete(key)

    async def get_or_set(self, key: str, factory, expire: int = 3600) -> Any:
        """Get từ cache, nếu không có thì gọi factory và cache."""
        value = await self.get(key)
        if value is not None:
            return value
        value = await factory()
        await self.set(key, value, expire)
        return value


# Sử dụng
@router.get("/users/{user_id}")
async def get_user(user_id: int, request: Request):
    cache = CacheService(request.app.state.redis)

    user = await cache.get_or_set(
        f"user:{user_id}",
        lambda: fetch_user_from_db(user_id),
        expire=300,
    )
    return user

3. 非同步深入研究

並發任務

import asyncio
import httpx
from fastapi import APIRouter

router = APIRouter()


@router.get("/dashboard")
async def dashboard():
    """Gọi nhiều services song song."""
    async with httpx.AsyncClient() as client:
        # Chạy 3 requests ĐỒNG THỜI
        users_task = client.get("http://user-service/users/count")
        posts_task = client.get("http://post-service/posts/count")
        orders_task = client.get("http://order-service/orders/recent")

        users_resp, posts_resp, orders_resp = await asyncio.gather(
            users_task, posts_task, orders_task,
            return_exceptions=True,  # Không crash nếu 1 task fail
        )

    return {
        "users": users_resp.json() if not isinstance(users_resp, Exception) else None,
        "posts": posts_resp.json() if not isinstance(posts_resp, Exception) else None,
        "orders": orders_resp.json() if not isinstance(orders_resp, Exception) else None,
    }

信號量 - 並發限制

import asyncio
import httpx

# Giới hạn tối đa 10 requests đồng thời
semaphore = asyncio.Semaphore(10)


async def fetch_with_limit(client: httpx.AsyncClient, url: str) -> dict:
    async with semaphore:
        response = await client.get(url)
        return response.json()


@router.post("/fetch-bulk")
async def fetch_bulk(urls: list[str]):
    """Fetch nhiều URLs với concurrency limit."""
    async with httpx.AsyncClient(timeout=30) as client:
        tasks = [fetch_with_limit(client, url) for url in urls[:100]]
        results = await asyncio.gather(*tasks, return_exceptions=True)

    return [
        r if not isinstance(r, Exception) else {"error": str(r)}
        for r in results
    ]

用於串流媒體的非同步生成器

from fastapi.responses import StreamingResponse


async def generate_csv(query_params: dict):
    """Stream CSV data - không load toàn bộ vào memory."""
    # Header
    yield "id,name,email,created_at\n"

    # Data rows - lấy từng batch
    offset = 0
    batch_size = 1000
    while True:
        users = await fetch_users_batch(offset, batch_size)
        if not users:
            break
        for user in users:
            yield f"{user.id},{user.name},{user.email},{user.created_at}\n"
        offset += batch_size


@router.get("/export/users")
async def export_users():
    """Export users ra CSV file (streaming)."""
    return StreamingResponse(
        generate_csv({}),
        media_type="text/csv",
        headers={"Content-Disposition": "attachment; filename=users.csv"},
    )

4. httpx - 非同步 HTTP 用戶端

# app/core/http_client.py
import httpx
from contextlib import asynccontextmanager


class HTTPClient:
    """Reusable async HTTP client."""

    def __init__(self):
        self.client: httpx.AsyncClient | None = None

    async def start(self):
        self.client = httpx.AsyncClient(
            timeout=httpx.Timeout(30.0, connect=10.0),
            limits=httpx.Limits(
                max_connections=100,
                max_keepalive_connections=20,
            ),
            follow_redirects=True,
        )

    async def stop(self):
        if self.client:
            await self.client.aclose()

    async def get(self, url: str, **kwargs) -> httpx.Response:
        if not self.client:
            raise RuntimeError("HTTP client not initialized")
        return await self.client.get(url, **kwargs)

    async def post(self, url: str, **kwargs) -> httpx.Response:
        if not self.client:
            raise RuntimeError("HTTP client not initialized")
        return await self.client.post(url, **kwargs)


http_client = HTTPClient()

# Trong lifespan
# await http_client.start()  # startup
# await http_client.stop()   # shutdown

總結

在這篇文章中,我們深入研究:

  • 文件上傳:單張、多張、串流上傳
  • Redis快取: fastapi-cache2 和自訂快取服務
  • 非同步深入研究:收集、信號量、非同步產生器
  • httpx:可重複使用的非同步 HTTP 用戶端

下一篇文章將討論大型應用程式的乾淨架構和專案結構。