FastAPIとAsync SQLAlchemy 2.0による非同期REST API構築:Connection Pool最適化とAlembicマイグレーション

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

1. 背景と実際の課題

本番環境で思わぬトラブルに遭遇したことがあります。FastAPIで丁寧に構築したAPIだったにもかかわらず、リクエスト数が1,500 req/sに達した途端、レイテンシが25msから800ms以上へと急増したのです。原因はよくあるものでした。従来の同期(sync)データベースドライバーによって、Uvicornの全ワーカーがブロックされてしまっていたのです。

FastAPIは、リクエスト受信からデータベースアクセスに至るI/Oパイプライン全体がノンブロッキングで実行されて初めて真のパフォーマンスを発揮します。SQLAlchemy 2.0は、asyncpgドライバーを通じたネイティブasync対応とselect()構文の標準化により、この課題を根本から解決してくれます。

しかし、非同期化への移行は単にawaitキーワードを付与するだけでは済みません。次の2つの大きな落とし穴に直面することになります。

  • Lazy-loading(遅延読み込み)の無効化: 非同期環境ではリレーションの暗黙的な読み込み機構が完全に無効化されます。不適切なコンテキストで呼び出すと、即座にMissingGreenletエラーが発生します。
  • Connection Poolの枯渇: プールサイズの設定を誤ると接続が滞留し、トラフィック急増時にAPIダウンを引き起こす原因になります。

2. インストールとディレクトリ構成

Python 3.10以降で仮想環境を作成し、必要なライブラリをインストールします。

# 仮想環境の作成
python -m venv .venv
source .venv/bin/activate

# FastAPI、asyncpgドライバー、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

明確にモジュール化されたプロジェクト構成:

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. システムのセットアップ

3.1. Async EngineとConnection Poolの設定

app/core/database.pyで、非同期エンジンと本番環境向けの主要なプールパラメータを初期化します。

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,          # プール内に常時維持する接続数(20件)
    max_overflow=10,       # 過負荷時に追加を許可する最大接続数(10件)
    pool_timeout=30,       # 接続取得のタイムアウト(30秒)
    pool_recycle=1800,     # 30分ごとに接続を再生成し、サイレント切断を防止
    pool_pre_ping=True,    # セッションへ渡す前にPingで生存確認を実施
)

AsyncSessionLocal = async_sessionmaker(
    bind=engine,
    class_=AsyncSession,
    autoflush=False,
    expire_on_commit=False, # commit後の不要なlazy-loadによるエラーを防止
)

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. SQLAlchemy 2.0準拠のモデル定義

app/models/item.pyで、型推論(type hinting)の恩恵を最大限に受けるためMappedおよびmapped_column構文を使用します。

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. Alembicによる非同期マイグレーションの統合

Alembicの設定ディレクトリを初期化します。

alembic init alembic

デフォルトのalembic/env.pyは同期エンジンのみをサポートしています。run_migrations_online関数を更新し、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  # モデルのメタデータを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. APIエンドポイントの実装

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()  # トランザクションを閉じずにSQLを発行して即座にIDを取得
    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. テストと実践的な運用

初回マイグレーションの作成と適用:

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

Uvicornサーバーの起動:

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

本番環境でのコネクションリークを監視するため、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(),
    }

運用時の重要ポイント:

  • セッションの即時解放: セッションブロック内で外部APIの呼び出しや負荷の高いCPU処理を実行することは絶対に避けてください。必要なデータを取得したら即座に接続を解放し、その後に後続処理を行います。
  • ワーカー数に応じたプール計算: Uvicornを4ワーカー(-w 4)で実行する場合、4つの独立したプールが生成されます。pool_size=20およびmax_overflow=10の場合、システム全体で最大4 x 30 = 120の接続を消費する可能性があります。PostgreSQLのmax_connectionsは必ずこの値より大きく設定してください。
  • pool_pre_ping=Trueの有効化: AWS RDSやネットワークプロキシによって切断されたステールな接続を自動的に検知・破棄し、ConnectionDoesNotExistErrorの発生を最小限に抑えます。
Share: