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

レッスン 16: ファイルのアップロード、キャッシュ、非同期の詳細

ファイルアップロードストリーミング、UploadFile、マルチパートフォーム。 fastapi-cache2 を使用した Redis キャッシュ。非同期/待機の詳細: 非同期、同時タスク、セマフォ、非同期ジェネレーター。非同期 HTTP クライアントの場合は httpx。

💻 プログラミング — レッスン 16 レッスン 16: ファイルのアップロード、キャッシュ、非同期の詳細 ダイブ

Python FastAPI: 基本から高度まで

パート 4: 高度な機能

xdev.asia

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 クライアント

次の記事では、大規模アプリのクリーン アーキテクチャとプロジェクト構造について説明します。