404 Not Found

404 Not Found


nginx

Database Integration — Asynchronous SQLAlchemy + Alembic

データベースは倉庫のようなものです — ドア (接続)が少なすぎると荷物 (クエリ)が詰まり, レイアウト (インデックス)が乱雑だと1つのアイテム (データ)を見つけるのに倉庫全体を探さなければなりません。

1. 学ぶ内容


2. Alice のリアルストーリー

(1) ペインポイント:データベースの同期処理が非同期 API を遅くする

Alice は同期 SQLAlchemy を使ってデータベースとやり取りしており, 各クエリがイベントループを 50-100 ms ブロックしています。PriceTracker の同時リクエスト数が 500 に達したとき, 非同期 FastAPI の利点が同期データベース操作によって完全に相殺され, P99 レイテンシが予想の 50 ms から 2,000 ms に急増しました。Charlie はデータベース接続プールが枯渇し, 新しいリクエストが接続を待ってキューに並んでいると指摘しました。

(2) 非同期 SQLAlchemy の解決策

SQLAlchemy 2.0 はネイティブの非同期サポートを提供します:create_async_engine + AsyncSession。データベースクエリがイベントループをブロックしなくなり, FastAPI が非同期の利点をフルに活用できるようになります。

PYTHON
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker

engine = create_async_engine("postgresql+asyncpg://user:pass@localhost/db")
async_session = async_sessionmaker(engine, expire_on_commit=False)

async def get_db():
    async with async_session() as session:
        yield session

(3) 成果

非同期 SQLAlchemy への移行後, PriceTracker の P99 レイテンシは 2,000 ms から 80 ms に低下し, データベース接続利用率は 30% から 90% に向上し, 単一ノードの QPS は 500 から 3,000 以上に上昇しました。


3. 非同期エンジンとセッション

(1) データベースアーキテクチャ ER 図

100%
erDiagram
    users ||--o{ products : creates
    products ||--o{ prices : has
    users {
        int id PK
        string email UK
        string hashed_password
        string role
        string subscription
        datetime created_at
    }
    products {
        int id PK
        string name
        string category
        float base_price
        string description
        int user_id FK
        datetime created_at
    }
    prices {
        int id PK
        int product_id FK
        float price
        string currency
        string source
        datetime recorded_at
    }

(1) ▶ サンプル:非同期エンジンと接続プール設定

PYTHON
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker

DATABASE_URL = "postgresql+asyncpg://pricetracker:secret@localhost:5432/pricetracker"

engine = create_async_engine(
    DATABASE_URL,
    echo=False,          # 開発時は True に設定して SQL ログを出力
    pool_size=20,        # 持続接続数
    max_overflow=10,     # プール枯渇時の追加接続
    pool_timeout=30,     # 利用可能接続の待ち時間
    pool_recycle=3600,   # 1時間後に接続をリサイクル
)

async_session = async_sessionmaker(
    engine,
    class_=AsyncSession,
    expire_on_commit=False,  # コミット後もオブジェクトにアクセス可能
)

出力:

TEXT
# 実行成功

(2) 同期 vs 非同期 SQLAlchemy の比較

側面 同期 SQLAlchemy 非同期 SQLAlchemy 2.0
エンジン create_engine create_async_engine
セッション Session AsyncSession
検索 session.execute(stmt) await session.execute(stmt)
コミット session.commit() await session.commit()
ドライバ psycopg2 asyncpg
ブロック あり なし
接続プール QueuePool AsyncAdaptedQueuePool

4. ORM モデル宣言 (新しい 2.0 スタイル)

(1) DeclarativeBase + Mapped + mapped_column

SQLAlchemy 2.0 は古い Column() 宣言を Mapped[type]mapped_column() に置き換え, 型ヒントと ORM マッピングを統一します。

(1) ▶ サンプル:PriceTracker 3テーブルモデル

PYTHON
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
from sqlalchemy import String, Float, Integer, DateTime, ForeignKey, Index
from sqlalchemy.orm import relationship
from datetime import datetime

class Base(DeclarativeBase):
    pass

class User(Base):
    __tablename__ = "users"

    id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
    email: Mapped[str] = mapped_column(String(255), unique=True, index=True)
    hashed_password: Mapped[str] = mapped_column(String(255))
    role: Mapped[str] = mapped_column(String(50), default="user")
    subscription: Mapped[str] = mapped_column(String(50), default="free")
    created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)

    products: Mapped[list["Product"]] = relationship(back_populates="owner")

class Product(Base):
    __tablename__ = "products"
    __table_args__ = (
        Index("ix_products_category", "category"),
        Index("ix_products_name", "name"),
    )

    id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
    name: Mapped[str] = mapped_column(String(200), nullable=False)
    category: Mapped[str] = mapped_column(String(100), nullable=False)
    base_price: Mapped[float] = mapped_column(Float, nullable=False)
    description: Mapped[str | None] = mapped_column(String(2000), nullable=True)
    user_id: Mapped[int] = mapped_column(Integer, ForeignKey("users.id"))
    created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)

    owner: Mapped["User"] = relationship(back_populates="products")
    prices: Mapped[list["Price"]] = relationship(back_populates="product")

class Price(Base):
    __tablename__ = "prices"
    __table_args__ = (
        Index("ix_prices_product_id", "product_id"),
        Index("ix_prices_recorded_at", "recorded_at"),
    )

    id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
    product_id: Mapped[int] = mapped_column(Integer, ForeignKey("products.id"))
    price: Mapped[float] = mapped_column(Float, nullable=False)
    currency: Mapped[str] = mapped_column(String(3), default="USD")
    source: Mapped[str] = mapped_column(String(100))
    recorded_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)

    product: Mapped["Product"] = relationship(back_populates="prices")

出力:

TEXT
# 実行成功

(2) V1 旧スタイル vs V2 新スタイル

側面 V1 旧スタイル V2 新スタイル
Base declarative_base() class Base(DeclarativeBase)
フィールド宣言 Column(Integer, primary_key=True) Mapped[int] = mapped_column(...)
型ヒント なし Mapped[type] 完全な型
オプショナルフィールド Column(String, nullable=True) Mapped[str | None]
リレーション relationship() Mapped[list["X"]] = relationship()

5. 非同期セッション管理と依存性注入

(1) 「yield」依存パターン

(1) ▶ サンプル:DB セッション依存

PYTHON
from typing import AsyncGenerator
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker

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

出力:

TEXT
# 関数定義成功

(2) ▶ サンプル:エンドポイントで非同期セッションを使用

PYTHON
from fastapi import FastAPI, Depends
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession

app = FastAPI()

@app.get("/products/{product_id}")
async def get_product(
    product_id: int,
    db: AsyncSession = Depends(get_db),
):
    stmt = select(Product).where(Product.id == product_id)
    result = await db.execute(stmt)
    product = result.scalar_one_or_none()
    if not product:
        from fastapi import HTTPException
        raise HTTPException(status_code=404, detail="Product not found")
    return {
        "id": product.id,
        "name": product.name,
        "base_price": product.base_price,
    }

出力:

TEXT
# 関数定義成功

6. Alembic 非同期マイグレーション

(1) マイグレーションワークフロー

100%
flowchart LR
    A[alembic revision --autogenerate -m desc] --> B[マイグレーションファイルの編集]
    B --> C[alembic upgrade head]
    C --> D[データベースに適用]
    D --> E{ロールバックが必要?}
    E -->|はい| F[alembic downgrade -1]
    E -->|いいえ| G[開発を継続]

(1) ▶ サンプル:Alembic の初期化

BASH
# Alembic のインストール
uv add alembic

# 非同期テンプレートで Alembic を初期化
cd pricetracker
alembic init -t async alembic

出力:

TEXT
# コマンド実行成功

(2) ▶ サンプル:alembic/env.py で非同期モードを設定

PYTHON
# alembic/env.py (主要セクション)
from sqlalchemy.ext.asyncio import create_async_engine
from app.models import Base  # モデルをインポート
from app.core.config import settings

target_metadata = Base.metadata

def run_migrations_online():
    connectable = create_async_engine(settings.database_url)

    async def do_run_migrations(connection):
        context = MigrationContext.configure(
            connection=connection,
            target_metadata=target_metadata,
        )
        with context.begin_transaction():
            context.run_migrations()

    with connectable.connect() as connection:
        asyncio.run(do_run_migrations(connection))

出力:

TEXT
# 関数定義成功

(3) ▶ サンプル:マイグレーションの生成と適用

BASH
# モデル変更からマイグレーションを自動生成
alembic revision --autogenerate -m "add users products prices tables"

# マイグレーションを適用
alembic upgrade head

# 1ステップ戻す
alembic downgrade -1

# 現在のバージョンを確認
alembic current

出力:

TEXT
# コマンド実行成功

❓ よくある質問

Q asyncpg と psycopg の違いは何ですか?
A asyncpg は純粋な非同期 PostgreSQL ドライバで, パフォーマンスに優れています。psycopg3 は非同期操作をサポートしていますが, C ライブラリベースです。本番環境では asyncpg を推奨します (postgresql+asyncpg://)。
Q expire_on_commit=False は何を意味しますか?
A デフォルトでは, コミット後に ORM オブジェクトのプロパティが期限切れになり, 再度アクセスするとクエリがトリガーされます。これを False に設定すると, プロパティに引き続きアクセスでき, 遅延読み込みの問題を防ぎます。
Q Alembic の autogenerate はすべての変更を検出しますか?
A いいえ。テーブルやカラムの追加・削除, インデックスの変更は検出しますが, カラム名の変更, 制約セマンティクスの変更などは検出しません。これらはマイグレーションファイルの手動編集が必要です。
Q 接続プールサイズはどのように設定すべきですか?
A 式:pool_size = (CPU コア数 x 2) + アクティブディスク数。PriceTracker は pool_size=20 + max_overflow=10 で数百万クエリを処理しています。
Q Mapped[str | None]Mapped[Optional[str]] に違いはありますか?
A 機能的には同等です。str | None は Python 3.10+ の構文, Optional[str]typing 互換の記法です。前者を推奨します。
Q 非同期コンテキストで同期 SQLAlchemy コードは使えますか?
A async 関数内で同期 SQLAlchemy メソッドを直接呼び出すことはできません。run_in_executor でラップするか, 非同期 API のみを使用してください。

📖 まとめ


📝 練習問題

  1. 基本問題 (難易度 ⭐):create_async_engine を設定して PostgreSQL に接続し, async_sessionmaker を作成し, get_db yield 依存を記述してください。ヒント:create_async_engine(DATABASE_URL, pool_size=20)
  2. 応用問題 (難易度 ⭐⭐):PriceTracker の ProductPrice モデルを定義してください (DeclarativeBase + Mapped スタイルを使用)。外部キーリレーションとインデックスを含めてください。ヒント:ForeignKey("products.id") + relationship()
  3. チャレンジ問題 (難易度 ⭐⭐⭐):Alembic を非同期モードで初期化し, env.py を設定し, autogenerate で3テーブルのマイグレーションファイルを生成し, upgrade head でデータベースに適用してください。ヒント:alembic init -t async alembic + env.pytarget_metadata を変更

---|

Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%