Database Integration — Asynchronous SQLAlchemy + Alembic
データベースは倉庫のようなものです — ドア (接続)が少なすぎると荷物 (クエリ)が詰まり, レイアウト (インデックス)が乱雑だと1つのアイテム (データ)を見つけるのに倉庫全体を探さなければなりません。
1. 学ぶ内容
- SQLAlchemy 2.0 非同期エンジン:
create_async_engine,AsyncSession - ORM モデル宣言:
DeclarativeBase,Mapped,mapped_column(新しい 2.0 スタイル) - Alembic 非同期マイグレーション:
env.pyでrun_migrations_online非同期モードを設定 - 非同期セッション管理:
async with/yield依存性注入パターン - Alice シナリオ:PriceTracker の
products/prices/users3テーブルの非同期モデル設計
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 図
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) マイグレーションワークフロー
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 のみを使用してください。📖 まとめ
- SQLAlchemy 2.0 非同期エンジン (
create_async_engine+AsyncSession)はイベントループをブロックせず, FastAPI の非同期の利点をフルに活用します Mapped[type]+mapped_column()は新しい 2.0 スタイルで, 型ヒントと ORM マッピングを統一しますyield依存パターンは非同期セッションのライフサイクルを管理し, 自動コミット, ロールバック, クローズを行います- Alembic 非同期マイグレーションは
-t asyncテンプレートで初期化し,env.pyテンプレートで非同期エンジンを設定します - PriceTracker の3テーブルモデル (users/products/prices)はインデックス戦略を含み, 数百万レコードのクエリをサポートします
📝 練習問題
- 基本問題 (難易度 ⭐):
create_async_engineを設定して PostgreSQL に接続し,async_sessionmakerを作成し,get_dbyield 依存を記述してください。ヒント:create_async_engine(DATABASE_URL, pool_size=20) - 応用問題 (難易度 ⭐⭐):PriceTracker の
ProductとPriceモデルを定義してください (DeclarativeBase + Mapped スタイルを使用)。外部キーリレーションとインデックスを含めてください。ヒント:ForeignKey("products.id")+relationship() - チャレンジ問題 (難易度 ⭐⭐⭐):Alembic を非同期モードで初期化し,
env.pyを設定し,autogenerateで3テーブルのマイグレーションファイルを生成し,upgrade headでデータベースに適用してください。ヒント:alembic init -t async alembic+env.pyのtarget_metadataを変更
---|



