Claude Skill

python-database-ops

SQLAlchemy and database patterns for Python. Triggers on: sqlalchemy, database, orm, migration, alembic, async database, connection pool, repository pattern, unit of work.

LLM Mart · 0 points · 0 views 0 listing impressions 0 install-command copies
Virus-scanned Reviewed automatically before listing.

Full trust report

Download 0xdarkmatter-claude-mods-skills_python-database-ops-3dfaf0b.zip · 12 KB
Part of 0xdarkmatter/claude-mods — 94 skills

Install

skills CLI npx skills add https://github.com/0xDarkMatter/claude-mods/tree/main/skills/python-database-ops
Claude Code claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install 0xdarkmatter-claude-mods@llmmart
Git git clone https://github.com/0xDarkMatter/claude-mods.git

The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole 0xdarkmatter/claude-mods collection as a plugin from our marketplace. Git is the plain clone.

Skill manifest

Python Database Patterns

SQLAlchemy 2.0 and database best practices.

SQLAlchemy 2.0 Basics

from sqlalchemy import create_engine, select
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, Session

class Base(DeclarativeBase):
    pass

class User(Base):
    __tablename__ = "users"

    id: Mapped[int] = mapped_column(primary_key=True)
    name: Mapped[str] = mapped_column(String(100))
    email: Mapped[str] = mapped_column(String(255), unique=True)
    is_active: Mapped[bool] = mapped_column(default=True)

# Create engine and tables
engine = create_engine("postgresql://user:pass@localhost/db")
Base.metadata.create_all(engine)

# Query with 2.0 style
with Session(engine) as session:
    stmt = select(User).where(User.is_active == True)
    users = session.execute(stmt).scalars().all()

Async SQLAlchemy

from sqlalchemy.ext.asyncio import (
    AsyncSession,
    async_sessionmaker,
    create_async_engine,
)
from sqlalchemy import select

# Async engine
engine = create_async_engine(
    "postgresql+asyncpg://user:pass@localhost/db",
    echo=False,
    pool_size=5,
    max_overflow=10,
)

# Session factory
async_session = async_sessionmaker(engine, expire_on_commit=False)

# Usage
async with async_session() as session:
    result = await session.execute(select(User).where(User.id == 1))
    user = result.scalar_one_or_none()

Model Relationships

from sqlalchemy import ForeignKey
from sqlalchemy.orm import relationship, Mapped, mapped_column

class User(Base):
    __tablename__ = "users"

    id: Mapped[int] = mapped_column(primary_key=True)
    name: Mapped[str]

    # One-to-many
    posts: Mapped[list["Post"]] = relationship(back_populates="author")

class Post(Base):
    __tablename__ = "posts"

    id: Mapped[int] = mapped_column(primary_key=True)
    title: Mapped[str]
    author_id: Mapped[int] = mapped_column(ForeignKey("users.id"))

    # Many-to-one
    author: Mapped["User"] = relationship(back_populates="posts")

Common Query Patterns

from sqlalchemy import select, and_, or_, func

# Basic select
stmt = select(User).where(User.is_active == True)

# Multiple conditions
stmt = select(User).where(
    and_(
        User.is_active == True,
        User.age >= 18
    )
)

# OR conditions
stmt = select(User).where(
    or_(User.role == "admin", User.role == "moderator")
)

# Ordering and limiting
stmt = select(User).order_by(User.created_at.desc()).limit(10)

# Aggregates
stmt = select(func.count(User.id)).where(User.is_active == True)

# Joins
stmt = select(User, Post).join(Post, User.id == Post.author_id)

# Eager loading
from sqlalchemy.orm import selectinload
stmt = select(User).options(selectinload(User.posts))

FastAPI Integration

from fastapi import Depends, FastAPI
from sqlalchemy.ext.asyncio import AsyncSession
from typing import Annotated

async def get_db() -> AsyncGenerator[AsyncSession, None]:
    async with async_session() as session:
        yield session

DB = Annotated[AsyncSession, Depends(get_db)]

@app.get("/users/{user_id}")
async def get_user(user_id: int, db: DB):
    result = await db.execute(select(User).where(User.id == user_id))
    user = result.scalar_one_or_none()
    if not user:
        raise HTTPException(status_code=404)
    return user

Quick Reference

Operation SQLAlchemy 2.0 Style
Select all select(User)
Filter .where(User.id == 1)
First .scalar_one_or_none()
All .scalars().all()
Count select(func.count(User.id))
Join .join(Post)
Eager load .options(selectinload(User.posts))

Additional Resources

  • ./references/sqlalchemy-async.md - Async patterns, session management
  • ./references/connection-pooling.md - Pool configuration, health checks
  • ./references/transactions.md - Transaction patterns, isolation levels
  • ./references/migrations.md - Alembic setup, migration strategies

Assets

  • ./assets/alembic.ini.template - Alembic configuration template

See Also

Prerequisites:

  • python-typing-ops - Mapped types and annotations
  • python-async-ops - Async database sessions

Related Skills:

  • python-fastapi-ops - Dependency injection for DB sessions
  • python-pytest-ops - Database fixtures and testing
Files (claude-mods)
  • assets
    • alembic.ini.template 1 KB · in bundle
  • references
    • connection-pooling.md 7.7 KB
      # Connection Pool Configuration
      
      Database connection pool patterns for production.
      
      ## SQLAlchemy Pool Settings
      
      ```python
      from sqlalchemy import create_engine
      from sqlalchemy.ext.asyncio import create_async_engine
      
      # Sync engine with pool config
      engine = create_engine(
          "postgresql://user:pass@localhost/db",
      
          # Pool size
          pool_size=5,           # Persistent connections (default: 5)
          max_overflow=10,       # Extra connections when pool exhausted
          # Total max connections = pool_size + max_overflow = 15
      
          # Timeouts
          pool_timeout=30,       # Wait for connection (seconds)
          pool_recycle=3600,     # Recycle connections after N seconds
          pool_pre_ping=True,    # Test connections before use
      
          # Connection args
          connect_args={
              "connect_timeout": 10,
              "options": "-c statement_timeout=30000",  # 30s query timeout
          },
      )
      
      
      # Async engine
      async_engine = create_async_engine(
          "postgresql+asyncpg://user:pass@localhost/db",
          pool_size=5,
          max_overflow=10,
          pool_timeout=30,
          pool_recycle=3600,
          pool_pre_ping=True,
      )
      ```
      
      ## Pool Sizing Guidelines
      
      ```python
      """
      Connection Pool Sizing
      
      Rule of thumb:
          pool_size = (CPU cores × 2) + disk spindles
      
      For async applications:
          pool_size = expected_concurrent_requests / avg_queries_per_request
      
      Examples:
          - Web app, 4 cores, SSD:  pool_size=10, max_overflow=10
          - Worker, 4 cores, HDD:   pool_size=12, max_overflow=5
          - High-traffic API:       pool_size=20, max_overflow=30
      """
      
      import os
      
      def calculate_pool_size() -> tuple[int, int]:
          """Calculate pool size based on environment."""
          cpu_count = os.cpu_count() or 4
      
          if os.getenv("ENV") == "production":
              pool_size = cpu_count * 2 + 4
              max_overflow = pool_size
          else:
              pool_size = 5
              max_overflow = 5
      
          return pool_size, max_overflow
      
      pool_size, max_overflow = calculate_pool_size()
      ```
      
      ## Pool Events and Monitoring
      
      ```python
      from sqlalchemy import event
      from sqlalchemy.pool import Pool
      import logging
      
      logger = logging.getLogger(__name__)
      
      @event.listens_for(Pool, "connect")
      def on_connect(dbapi_conn, connection_record):
          """Called when a new connection is created."""
          logger.debug("New database connection created")
      
      @event.listens_for(Pool, "checkout")
      def on_checkout(dbapi_conn, connection_record, connection_proxy):
          """Called when a connection is retrieved from pool."""
          logger.debug("Connection checked out from pool")
      
      @event.listens_for(Pool, "checkin")
      def on_checkin(dbapi_conn, connection_record):
          """Called when a connection is returned to pool."""
          logger.debug("Connection returned to pool")
      
      @event.listens_for(Pool, "invalidate")
      def on_invalidate(dbapi_conn, connection_record, exception):
          """Called when a connection is invalidated."""
          logger.warning(f"Connection invalidated: {exception}")
      
      
      # Pool statistics
      def log_pool_status(engine):
          """Log current pool status."""
          pool = engine.pool
          logger.info(
              f"Pool status: "
              f"size={pool.size()}, "
              f"checked_out={pool.checkedout()}, "
              f"overflow={pool.overflow()}, "
              f"checkedin={pool.checkedin()}"
          )
      ```
      
      ## Health Check Endpoint
      
      ```python
      from fastapi import FastAPI, HTTPException
      from sqlalchemy import text
      import asyncio
      
      app = FastAPI()
      
      async def check_database_health(timeout: float = 5.0) -> dict:
          """Check database connectivity and response time."""
          try:
              start = asyncio.get_event_loop().time()
      
              async with async_session_factory() as session:
                  await asyncio.wait_for(
                      session.execute(text("SELECT 1")),
                      timeout=timeout
                  )
      
              latency = (asyncio.get_event_loop().time() - start) * 1000
      
              return {
                  "status": "healthy",
                  "latency_ms": round(latency, 2),
                  "pool_size": async_engine.pool.size(),
                  "pool_checked_out": async_engine.pool.checkedout(),
              }
          except asyncio.TimeoutError:
              return {"status": "unhealthy", "error": "timeout"}
          except Exception as e:
              return {"status": "unhealthy", "error": str(e)}
      
      
      @app.get("/health/db")
      async def database_health():
          health = await check_database_health()
          if health["status"] != "healthy":
              raise HTTPException(status_code=503, detail=health)
          return health
      ```
      
      ## Connection Pool per Service
      
      ```python
      from dataclasses import dataclass
      from sqlalchemy.ext.asyncio import AsyncEngine, create_async_engine
      
      @dataclass
      class DatabaseConfig:
          url: str
          pool_size: int = 5
          max_overflow: int = 10
          pool_timeout: int = 30
          pool_recycle: int = 3600
      
      class DatabasePool:
          """Manage multiple database connections."""
      
          def __init__(self):
              self._engines: dict[str, AsyncEngine] = {}
      
          def add_database(self, name: str, config: DatabaseConfig):
              """Add a database connection pool."""
              self._engines[name] = create_async_engine(
                  config.url,
                  pool_size=config.pool_size,
                  max_overflow=config.max_overflow,
                  pool_timeout=config.pool_timeout,
                  pool_recycle=config.pool_recycle,
                  pool_pre_ping=True,
              )
      
          def get_engine(self, name: str) -> AsyncEngine:
              return self._engines[name]
      
          async def close_all(self):
              """Close all connection pools."""
              for engine in self._engines.values():
                  await engine.dispose()
      
      
      # Usage
      db_pool = DatabasePool()
      
      db_pool.add_database("primary", DatabaseConfig(
          url="postgresql+asyncpg://user:pass@primary/db",
          pool_size=10,
      ))
      
      db_pool.add_database("replica", DatabaseConfig(
          url="postgresql+asyncpg://user:pass@replica/db",
          pool_size=20,  # More connections for read replica
      ))
      ```
      
      ## Read/Write Splitting
      
      ```python
      from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
      
      # Separate session factories for read/write
      write_engine = create_async_engine(
          "postgresql+asyncpg://user:pass@primary/db",
          pool_size=10,
      )
      
      read_engine = create_async_engine(
          "postgresql+asyncpg://user:pass@replica/db",
          pool_size=20,
      )
      
      write_session = async_sessionmaker(write_engine, expire_on_commit=False)
      read_session = async_sessionmaker(read_engine, expire_on_commit=False)
      
      
      # FastAPI dependencies
      async def get_write_db():
          async with write_session() as session:
              yield session
      
      async def get_read_db():
          async with read_session() as session:
              yield session
      
      WriteDB = Annotated[AsyncSession, Depends(get_write_db)]
      ReadDB = Annotated[AsyncSession, Depends(get_read_db)]
      
      
      @app.get("/users")
      async def list_users(db: ReadDB):  # Read from replica
          result = await db.execute(select(User))
          return result.scalars().all()
      
      @app.post("/users")
      async def create_user(user: UserCreate, db: WriteDB):  # Write to primary
          db_user = User(**user.model_dump())
          db.add(db_user)
          await db.commit()
          return db_user
      ```
      
      ## Graceful Shutdown
      
      ```python
      from contextlib import asynccontextmanager
      from fastapi import FastAPI
      
      @asynccontextmanager
      async def lifespan(app: FastAPI):
          # Startup - engines already created
          yield
          # Shutdown - close all pools gracefully
          await async_engine.dispose()
          logger.info("Database connections closed")
      
      app = FastAPI(lifespan=lifespan)
      ```
      
      ## Quick Reference
      
      | Setting | Purpose | Typical Value |
      |---------|---------|---------------|
      | `pool_size` | Persistent connections | 5-20 |
      | `max_overflow` | Extra connections | 10-30 |
      | `pool_timeout` | Wait for connection | 30s |
      | `pool_recycle` | Recycle connection age | 3600s |
      | `pool_pre_ping` | Test before use | True |
      
      | Scenario | pool_size | max_overflow |
      |----------|-----------|--------------|
      | Development | 5 | 5 |
      | Small API | 10 | 10 |
      | High-traffic | 20 | 30 |
      | Background worker | 5 | 5 |
      
    • migrations.md 7.7 KB
      # Database Migrations with Alembic
      
      Schema migration patterns for SQLAlchemy projects.
      
      ## Setup
      
      ```bash
      # Install (project dependency)
      uv add alembic
      
      # Initialize in project root
      uv run alembic init alembic
      
      # For async projects
      alembic init -t async alembic
      ```
      
      ## Configuration
      
      ```python
      # alembic/env.py
      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.models import Base  # Your declarative base
      from app.config import settings
      
      config = context.config
      
      # Set database URL from settings
      config.set_main_option("sqlalchemy.url", settings.database_url)
      
      target_metadata = Base.metadata
      
      
      def run_migrations_offline():
          """Run migrations in 'offline' mode."""
          url = config.get_main_option("sqlalchemy.url")
          context.configure(
              url=url,
              target_metadata=target_metadata,
              literal_binds=True,
              dialect_opts={"paramstyle": "named"},
          )
          with context.begin_transaction():
              context.run_migrations()
      
      
      def do_run_migrations(connection: Connection):
          context.configure(connection=connection, target_metadata=target_metadata)
          with context.begin_transaction():
              context.run_migrations()
      
      
      async def run_async_migrations():
          """Run migrations in 'online' mode with async engine."""
          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():
          import asyncio
          asyncio.run(run_async_migrations())
      
      
      if context.is_offline_mode():
          run_migrations_offline()
      else:
          run_migrations_online()
      ```
      
      ## Common Commands
      
      ```bash
      # Generate migration from model changes
      alembic revision --autogenerate -m "add users table"
      
      # Apply all pending migrations
      alembic upgrade head
      
      # Rollback one migration
      alembic downgrade -1
      
      # Rollback to specific revision
      alembic downgrade abc123
      
      # Show current revision
      alembic current
      
      # Show migration history
      alembic history
      
      # Show pending migrations
      alembic history --indicate-current
      ```
      
      ## Migration Script Example
      
      ```python
      """Add users table
      
      Revision ID: abc123
      Revises:
      Create Date: 2024-01-15 10:00:00.000000
      """
      from typing import Sequence
      from alembic import op
      import sqlalchemy as sa
      
      revision: str = 'abc123'
      down_revision: str | None = None
      branch_labels: str | Sequence[str] | None = None
      depends_on: str | Sequence[str] | None = None
      
      
      def upgrade() -> None:
          op.create_table(
              'users',
              sa.Column('id', sa.Integer(), primary_key=True),
              sa.Column('email', sa.String(255), nullable=False, unique=True),
              sa.Column('name', sa.String(100), nullable=False),
              sa.Column('is_active', sa.Boolean(), default=True),
              sa.Column('created_at', sa.DateTime(), server_default=sa.func.now()),
          )
          op.create_index('ix_users_email', 'users', ['email'])
      
      
      def downgrade() -> None:
          op.drop_index('ix_users_email')
          op.drop_table('users')
      ```
      
      ## Data Migrations
      
      ```python
      """Migrate user names to lowercase
      
      Revision ID: def456
      """
      from alembic import op
      import sqlalchemy as sa
      from sqlalchemy.sql import table, column
      
      revision = 'def456'
      down_revision = 'abc123'
      
      
      def upgrade() -> None:
          # Define table structure for data migration
          users = table(
              'users',
              column('id', sa.Integer),
              column('name', sa.String),
          )
      
          # Update data
          op.execute(
              users.update().values(name=sa.func.lower(users.c.name))
          )
      
      
      def downgrade() -> None:
          # Data migrations are often one-way
          pass
      
      
      # For complex data migrations
      def upgrade() -> None:
          connection = op.get_bind()
      
          # Read in batches
          results = connection.execute(
              sa.text("SELECT id, name FROM users")
          )
      
          for batch in results.partitions(1000):
              for row in batch:
                  connection.execute(
                      sa.text("UPDATE users SET name = :name WHERE id = :id"),
                      {"id": row.id, "name": row.name.lower()}
                  )
      ```
      
      ## Adding Columns Safely
      
      ```python
      """Add nullable column first, then populate
      
      Production-safe column addition for large tables.
      """
      
      def upgrade() -> None:
          # Step 1: Add nullable column (fast, no table rewrite)
          op.add_column(
              'users',
              sa.Column('phone', sa.String(20), nullable=True)
          )
      
          # Step 2: Populate data (can be done in batches)
          # This is often done in a separate migration or script
      
          # Step 3: Add constraint (in a later migration after data is populated)
          # op.alter_column('users', 'phone', nullable=False)
      
      
      def downgrade() -> None:
          op.drop_column('users', 'phone')
      ```
      
      ## Renaming Columns
      
      ```python
      """Rename column with zero downtime
      
      Use a multi-step approach for production.
      """
      
      # Migration 1: Add new column
      def upgrade() -> None:
          op.add_column('users', sa.Column('full_name', sa.String(200)))
          # Copy data
          op.execute("UPDATE users SET full_name = name")
      
      def downgrade() -> None:
          op.drop_column('users', 'full_name')
      
      
      # Migration 2: Drop old column (after app updated to use new column)
      def upgrade() -> None:
          op.drop_column('users', 'name')
      
      def downgrade() -> None:
          op.add_column('users', sa.Column('name', sa.String(100)))
          op.execute("UPDATE users SET name = full_name")
      ```
      
      ## Index Management
      
      ```python
      """Add index concurrently (PostgreSQL)
      
      Non-blocking index creation for large tables.
      """
      from alembic import op
      
      def upgrade() -> None:
          # Create index without locking table (PostgreSQL)
          op.execute("""
              CREATE INDEX CONCURRENTLY IF NOT EXISTS
              ix_users_created_at ON users (created_at)
          """)
      
      
      def downgrade() -> None:
          op.execute("DROP INDEX CONCURRENTLY IF EXISTS ix_users_created_at")
      
      
      # Note: CONCURRENTLY cannot run inside a transaction
      # Add to migration script:
      # from alembic import context
      # context.execute_ddl_statements = True
      ```
      
      ## Multi-Database Migrations
      
      ```python
      # alembic.ini
      [alembic]
      script_location = alembic
      
      [primary]
      sqlalchemy.url = postgresql://user:pass@primary/db
      
      [analytics]
      sqlalchemy.url = postgresql://user:pass@analytics/db
      ```
      
      ```bash
      # Run migrations for specific database
      alembic -n primary upgrade head
      alembic -n analytics upgrade head
      ```
      
      ## Testing Migrations
      
      ```python
      import pytest
      from alembic import command
      from alembic.config import Config
      
      @pytest.fixture
      def alembic_config():
          config = Config("alembic.ini")
          config.set_main_option("sqlalchemy.url", "sqlite:///:memory:")
          return config
      
      def test_migrations_up_down(alembic_config):
          """Test that all migrations apply and rollback cleanly."""
          # Apply all migrations
          command.upgrade(alembic_config, "head")
      
          # Rollback all migrations
          command.downgrade(alembic_config, "base")
      
          # Apply again
          command.upgrade(alembic_config, "head")
      
      
      def test_migration_idempotent(alembic_config):
          """Test migrations can be run multiple times."""
          command.upgrade(alembic_config, "head")
          command.upgrade(alembic_config, "head")  # Should be no-op
      ```
      
      ## Quick Reference
      
      | Command | Purpose |
      |---------|---------|
      | `alembic revision --autogenerate -m "msg"` | Generate migration |
      | `alembic upgrade head` | Apply all migrations |
      | `alembic downgrade -1` | Rollback one |
      | `alembic current` | Show current version |
      | `alembic history` | List all migrations |
      
      | Operation | Method |
      |-----------|--------|
      | Create table | `op.create_table()` |
      | Drop table | `op.drop_table()` |
      | Add column | `op.add_column()` |
      | Drop column | `op.drop_column()` |
      | Alter column | `op.alter_column()` |
      | Create index | `op.create_index()` |
      | Execute SQL | `op.execute()` |
      
    • sqlalchemy-async.md 8.2 KB
      # Async SQLAlchemy Patterns
      
      Modern async database patterns with SQLAlchemy 2.0.
      
      ## Engine and Session Setup
      
      ```python
      from sqlalchemy.ext.asyncio import (
          AsyncSession,
          AsyncEngine,
          async_sessionmaker,
          create_async_engine,
      )
      
      # Create async engine
      engine = create_async_engine(
          "postgresql+asyncpg://user:pass@localhost/db",
          echo=False,            # SQL logging
          pool_size=5,           # Connection pool size
          max_overflow=10,       # Extra connections allowed
          pool_pre_ping=True,    # Test connections before use
          pool_recycle=3600,     # Recycle connections after 1 hour
      )
      
      # Session factory (not the session itself)
      async_session_factory = async_sessionmaker(
          engine,
          class_=AsyncSession,
          expire_on_commit=False,  # Don't expire objects after commit
      )
      
      
      # Usage with context manager
      async def get_users():
          async with async_session_factory() as session:
              result = await session.execute(select(User))
              return result.scalars().all()
      ```
      
      ## Session Scopes
      
      ```python
      # Per-request (FastAPI dependency)
      async def get_db():
          async with async_session_factory() as session:
              try:
                  yield session
                  await session.commit()
              except Exception:
                  await session.rollback()
                  raise
      
      
      # Explicit transaction control
      async def transfer_funds(from_id: int, to_id: int, amount: Decimal):
          async with async_session_factory() as session:
              async with session.begin():  # Auto-commit on success
                  from_account = await session.get(Account, from_id)
                  to_account = await session.get(Account, to_id)
      
                  from_account.balance -= amount
                  to_account.balance += amount
                  # Commits automatically if no exception
      
      
      # Nested transactions (savepoints)
      async def complex_operation():
          async with async_session_factory() as session:
              async with session.begin():
                  # Outer transaction
                  user = User(name="Test")
                  session.add(user)
      
                  try:
                      async with session.begin_nested():  # Savepoint
                          # Inner operation that might fail
                          await risky_operation(session)
                  except RiskyOperationError:
                      # Savepoint rolled back, outer continues
                      pass
      
                  await session.commit()
      ```
      
      ## Lazy Loading in Async
      
      ```python
      from sqlalchemy.orm import selectinload, joinedload, subqueryload
      
      # WRONG - lazy loading doesn't work in async
      async def bad_example():
          async with async_session_factory() as session:
              user = await session.get(User, 1)
              # This raises an error!
              print(user.posts)  # MissingGreenlet error
      
      
      # CORRECT - eager loading
      async def good_example():
          async with async_session_factory() as session:
              # Option 1: selectinload (separate query per relationship)
              stmt = select(User).options(selectinload(User.posts))
              result = await session.execute(stmt)
              user = result.scalar_one()
              print(user.posts)  # Works!
      
              # Option 2: joinedload (single JOIN query)
              stmt = select(User).options(joinedload(User.profile))
              result = await session.execute(stmt)
              user = result.scalar_one()
      
      
      # With nested relationships
      stmt = select(User).options(
          selectinload(User.posts).selectinload(Post.comments)
      )
      ```
      
      ## Async Session Dependency
      
      ```python
      from fastapi import Depends
      from typing import Annotated, AsyncGenerator
      
      async def get_db() -> AsyncGenerator[AsyncSession, None]:
          """Dependency for FastAPI."""
          async with async_session_factory() as session:
              yield session
      
      DB = Annotated[AsyncSession, Depends(get_db)]
      
      
      # With automatic transaction handling
      async def get_db_with_transaction() -> AsyncGenerator[AsyncSession, None]:
          async with async_session_factory() as session:
              try:
                  yield session
                  await session.commit()
              except Exception:
                  await session.rollback()
                  raise
              finally:
                  await session.close()
      ```
      
      ## Batch Operations
      
      ```python
      from sqlalchemy import insert, update, delete
      
      # Bulk insert
      async def bulk_create_users(users_data: list[dict]):
          async with async_session_factory() as session:
              stmt = insert(User).values(users_data)
              await session.execute(stmt)
              await session.commit()
      
      
      # Bulk update
      async def deactivate_users(user_ids: list[int]):
          async with async_session_factory() as session:
              stmt = (
                  update(User)
                  .where(User.id.in_(user_ids))
                  .values(is_active=False)
              )
              result = await session.execute(stmt)
              await session.commit()
              return result.rowcount
      
      
      # Bulk delete
      async def delete_old_posts(before_date: datetime):
          async with async_session_factory() as session:
              stmt = delete(Post).where(Post.created_at < before_date)
              result = await session.execute(stmt)
              await session.commit()
              return result.rowcount
      
      
      # Batch processing with chunks
      async def process_all_users(batch_size: int = 100):
          async with async_session_factory() as session:
              offset = 0
              while True:
                  stmt = select(User).offset(offset).limit(batch_size)
                  result = await session.execute(stmt)
                  users = result.scalars().all()
      
                  if not users:
                      break
      
                  for user in users:
                      await process_user(user)
      
                  await session.commit()
                  offset += batch_size
      ```
      
      ## Streaming Results
      
      ```python
      from sqlalchemy import select
      
      async def stream_large_table():
          """Process large tables without loading all into memory."""
          async with async_session_factory() as session:
              stmt = select(User).execution_options(yield_per=100)
              result = await session.stream(stmt)
      
              async for user in result.scalars():
                  await process_user(user)
      
      
      # Partitioned streaming
      async def stream_partitioned():
          async with async_session_factory() as session:
              stmt = select(User).execution_options(yield_per=100)
              result = await session.stream(stmt)
      
              async for partition in result.scalars().partitions(100):
                  # partition is a list of 100 users
                  await process_batch(partition)
      ```
      
      ## Async Raw SQL
      
      ```python
      from sqlalchemy import text
      
      async def raw_query():
          async with async_session_factory() as session:
              # Simple query
              result = await session.execute(
                  text("SELECT * FROM users WHERE is_active = :active"),
                  {"active": True}
              )
              rows = result.fetchall()
      
              # With column access
              for row in rows:
                  print(row.id, row.name)
      
      
      async def raw_insert():
          async with async_session_factory() as session:
              await session.execute(
                  text("INSERT INTO logs (message) VALUES (:msg)"),
                  {"msg": "Test log"}
              )
              await session.commit()
      ```
      
      ## Testing with Async
      
      ```python
      import pytest
      import pytest_asyncio
      from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
      
      @pytest_asyncio.fixture(scope="session")
      async def async_engine():
          engine = create_async_engine("sqlite+aiosqlite:///:memory:")
          async with engine.begin() as conn:
              await conn.run_sync(Base.metadata.create_all)
          yield engine
          await engine.dispose()
      
      @pytest_asyncio.fixture
      async def async_session(async_engine):
          async with AsyncSession(async_engine) as session:
              async with session.begin():
                  yield session
                  await session.rollback()
      
      
      @pytest.mark.asyncio
      async def test_create_user(async_session):
          user = User(name="Test", email="test@example.com")
          async_session.add(user)
          await async_session.flush()
      
          assert user.id is not None
      ```
      
      ## Quick Reference
      
      | Pattern | Async SQLAlchemy |
      |---------|------------------|
      | Create engine | `create_async_engine(url)` |
      | Session factory | `async_sessionmaker(engine)` |
      | Get session | `async with factory() as session:` |
      | Execute | `await session.execute(stmt)` |
      | Get one | `result.scalar_one_or_none()` |
      | Get all | `result.scalars().all()` |
      | Stream | `await session.stream(stmt)` |
      | Commit | `await session.commit()` |
      | Transaction | `async with session.begin():` |
      | Eager load | `.options(selectinload(rel))` |
      
    • transactions.md 7.7 KB
      # Transaction Patterns
      
      Database transaction management for data integrity.
      
      ## Basic Transaction Patterns
      
      ```python
      from sqlalchemy.ext.asyncio import AsyncSession
      
      # Pattern 1: Context manager (auto-commit/rollback)
      async with async_session_factory() as session:
          async with session.begin():
              user = User(name="Test")
              session.add(user)
              # Auto-commits on exit, rollback on exception
      
      
      # Pattern 2: Explicit control
      async with async_session_factory() as session:
          try:
              user = User(name="Test")
              session.add(user)
              await session.commit()
          except Exception:
              await session.rollback()
              raise
      
      
      # Pattern 3: Dependency with auto-commit
      async def get_db():
          async with async_session_factory() as session:
              try:
                  yield session
                  await session.commit()
              except Exception:
                  await session.rollback()
                  raise
      ```
      
      ## Nested Transactions (Savepoints)
      
      ```python
      async def complex_operation():
          async with async_session_factory() as session:
              async with session.begin():
                  # Create user (outer transaction)
                  user = User(name="Test")
                  session.add(user)
                  await session.flush()  # Get user.id
      
                  try:
                      # Nested operation (savepoint)
                      async with session.begin_nested():
                          profile = Profile(user_id=user.id, bio="Hello")
                          session.add(profile)
                          await session.flush()
      
                          # This might fail
                          await validate_profile(profile)
      
                  except ValidationError:
                      # Savepoint rolled back, but user is preserved
                      logger.warning("Profile creation failed")
      
                  # Commit user (profile may or may not exist)
                  await session.commit()
      ```
      
      ## Unit of Work Pattern
      
      ```python
      from typing import TypeVar, Generic
      from sqlalchemy.ext.asyncio import AsyncSession
      
      T = TypeVar("T")
      
      class UnitOfWork:
          """Coordinate multiple repository operations in a transaction."""
      
          def __init__(self, session_factory):
              self._session_factory = session_factory
              self._session: AsyncSession | None = None
      
          async def __aenter__(self):
              self._session = self._session_factory()
              return self
      
          async def __aexit__(self, exc_type, exc_val, exc_tb):
              if exc_type:
                  await self.rollback()
              await self._session.close()
      
          async def commit(self):
              await self._session.commit()
      
          async def rollback(self):
              await self._session.rollback()
      
          @property
          def users(self) -> "UserRepository":
              return UserRepository(self._session)
      
          @property
          def orders(self) -> "OrderRepository":
              return OrderRepository(self._session)
      
      
      # Usage
      async def create_order_with_items(user_id: int, items: list):
          async with UnitOfWork(async_session_factory) as uow:
              user = await uow.users.get(user_id)
              if not user:
                  raise NotFoundError("User not found")
      
              order = Order(user_id=user_id)
              order = await uow.orders.add(order)
      
              for item in items:
                  await uow.orders.add_item(order.id, item)
      
              await uow.commit()
              return order
      ```
      
      ## Isolation Levels
      
      ```python
      from sqlalchemy import create_engine
      from sqlalchemy.orm import Session
      
      # Engine-level default
      engine = create_engine(
          "postgresql://...",
          isolation_level="REPEATABLE READ"  # Default for all sessions
      )
      
      
      # Per-session isolation
      async with async_session_factory() as session:
          await session.connection(
              execution_options={"isolation_level": "SERIALIZABLE"}
          )
          # This session uses SERIALIZABLE isolation
      
      
      # Transaction-level in raw SQL
      async with async_session_factory() as session:
          await session.execute(text("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE"))
          # Perform operations
          await session.commit()
      ```
      
      ### Isolation Level Reference
      
      | Level | Dirty Read | Non-repeatable Read | Phantom Read |
      |-------|------------|---------------------|--------------|
      | READ UNCOMMITTED | Possible | Possible | Possible |
      | READ COMMITTED | No | Possible | Possible |
      | REPEATABLE READ | No | No | Possible* |
      | SERIALIZABLE | No | No | No |
      
      *PostgreSQL prevents phantoms in REPEATABLE READ
      
      ## Optimistic Locking
      
      ```python
      from sqlalchemy import Column, Integer
      from sqlalchemy.orm import Mapped, mapped_column
      
      class Account(Base):
          __tablename__ = "accounts"
      
          id: Mapped[int] = mapped_column(primary_key=True)
          balance: Mapped[int]
          version: Mapped[int] = mapped_column(default=0)
      
          __mapper_args__ = {"version_id_col": version}
      
      
      async def transfer_funds(from_id: int, to_id: int, amount: int):
          """Transfer with optimistic locking."""
          async with async_session_factory() as session:
              from_account = await session.get(Account, from_id)
              to_account = await session.get(Account, to_id)
      
              if from_account.balance < amount:
                  raise InsufficientFunds()
      
              from_account.balance -= amount
              to_account.balance += amount
      
              try:
                  await session.commit()
              except StaleDataError:
                  # Concurrent modification detected
                  await session.rollback()
                  raise ConcurrentModificationError()
      ```
      
      ## Pessimistic Locking
      
      ```python
      from sqlalchemy import select
      
      async def transfer_with_lock(from_id: int, to_id: int, amount: int):
          """Transfer with row-level lock."""
          async with async_session_factory() as session:
              async with session.begin():
                  # Lock rows for update
                  stmt = (
                      select(Account)
                      .where(Account.id.in_([from_id, to_id]))
                      .with_for_update()  # SELECT ... FOR UPDATE
                  )
                  result = await session.execute(stmt)
                  accounts = {a.id: a for a in result.scalars()}
      
                  from_account = accounts[from_id]
                  to_account = accounts[to_id]
      
                  if from_account.balance < amount:
                      raise InsufficientFunds()
      
                  from_account.balance -= amount
                  to_account.balance += amount
                  # Commit releases locks
      
      
      # Lock with options
      stmt = select(Account).with_for_update(
          nowait=True,     # Fail immediately if locked
          skip_locked=True # Skip locked rows (for queue processing)
      )
      ```
      
      ## Retry on Serialization Failure
      
      ```python
      from sqlalchemy.exc import OperationalError
      import asyncio
      
      async def retry_on_conflict(
          func,
          max_retries: int = 3,
          base_delay: float = 0.1,
      ):
          """Retry transaction on serialization failure."""
          for attempt in range(max_retries):
              try:
                  return await func()
              except OperationalError as e:
                  if "serialization" in str(e).lower() or "deadlock" in str(e).lower():
                      if attempt < max_retries - 1:
                          delay = base_delay * (2 ** attempt)
                          await asyncio.sleep(delay)
                          continue
                  raise
      
      
      # Usage
      async def process_order(order_id: int):
          async def _process():
              async with async_session_factory() as session:
                  async with session.begin():
                      order = await session.get(Order, order_id)
                      order.status = "processed"
                      await session.commit()
      
          await retry_on_conflict(_process)
      ```
      
      ## Quick Reference
      
      | Pattern | Use Case |
      |---------|----------|
      | `session.begin()` | Auto-commit/rollback |
      | `session.begin_nested()` | Savepoint (partial rollback) |
      | `with_for_update()` | Row-level locking |
      | `version_id_col` | Optimistic concurrency |
      | Isolation levels | Control visibility |
      
      | Isolation Level | When to Use |
      |-----------------|-------------|
      | READ COMMITTED | Default, most apps |
      | REPEATABLE READ | Reports, analytics |
      | SERIALIZABLE | Financial, inventory |
      
  • scripts
    • .gitkeep 0 B · in bundle
  • SKILL.md 4.7 KB
    ---
    name: python-database-ops
    description: "SQLAlchemy and database patterns for Python. Triggers on: sqlalchemy, database, orm, migration, alembic, async database, connection pool, repository pattern, unit of work."
    license: MIT
    compatibility: "SQLAlchemy 2.0+, Python 3.10+. Async requires asyncpg (PostgreSQL) or aiosqlite."
    allowed-tools: "Read Write Bash"
    metadata:
      author: claude-mods
      depends-on: python-typing-ops, python-async-ops
      related-skills: python-fastapi-ops, postgres-ops
    ---
    
    # Python Database Patterns
    
    SQLAlchemy 2.0 and database best practices.
    
    ## SQLAlchemy 2.0 Basics
    
    ```python
    from sqlalchemy import create_engine, select
    from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, Session
    
    class Base(DeclarativeBase):
        pass
    
    class User(Base):
        __tablename__ = "users"
    
        id: Mapped[int] = mapped_column(primary_key=True)
        name: Mapped[str] = mapped_column(String(100))
        email: Mapped[str] = mapped_column(String(255), unique=True)
        is_active: Mapped[bool] = mapped_column(default=True)
    
    # Create engine and tables
    engine = create_engine("postgresql://user:pass@localhost/db")
    Base.metadata.create_all(engine)
    
    # Query with 2.0 style
    with Session(engine) as session:
        stmt = select(User).where(User.is_active == True)
        users = session.execute(stmt).scalars().all()
    ```
    
    ## Async SQLAlchemy
    
    ```python
    from sqlalchemy.ext.asyncio import (
        AsyncSession,
        async_sessionmaker,
        create_async_engine,
    )
    from sqlalchemy import select
    
    # Async engine
    engine = create_async_engine(
        "postgresql+asyncpg://user:pass@localhost/db",
        echo=False,
        pool_size=5,
        max_overflow=10,
    )
    
    # Session factory
    async_session = async_sessionmaker(engine, expire_on_commit=False)
    
    # Usage
    async with async_session() as session:
        result = await session.execute(select(User).where(User.id == 1))
        user = result.scalar_one_or_none()
    ```
    
    ## Model Relationships
    
    ```python
    from sqlalchemy import ForeignKey
    from sqlalchemy.orm import relationship, Mapped, mapped_column
    
    class User(Base):
        __tablename__ = "users"
    
        id: Mapped[int] = mapped_column(primary_key=True)
        name: Mapped[str]
    
        # One-to-many
        posts: Mapped[list["Post"]] = relationship(back_populates="author")
    
    class Post(Base):
        __tablename__ = "posts"
    
        id: Mapped[int] = mapped_column(primary_key=True)
        title: Mapped[str]
        author_id: Mapped[int] = mapped_column(ForeignKey("users.id"))
    
        # Many-to-one
        author: Mapped["User"] = relationship(back_populates="posts")
    ```
    
    ## Common Query Patterns
    
    ```python
    from sqlalchemy import select, and_, or_, func
    
    # Basic select
    stmt = select(User).where(User.is_active == True)
    
    # Multiple conditions
    stmt = select(User).where(
        and_(
            User.is_active == True,
            User.age >= 18
        )
    )
    
    # OR conditions
    stmt = select(User).where(
        or_(User.role == "admin", User.role == "moderator")
    )
    
    # Ordering and limiting
    stmt = select(User).order_by(User.created_at.desc()).limit(10)
    
    # Aggregates
    stmt = select(func.count(User.id)).where(User.is_active == True)
    
    # Joins
    stmt = select(User, Post).join(Post, User.id == Post.author_id)
    
    # Eager loading
    from sqlalchemy.orm import selectinload
    stmt = select(User).options(selectinload(User.posts))
    ```
    
    ## FastAPI Integration
    
    ```python
    from fastapi import Depends, FastAPI
    from sqlalchemy.ext.asyncio import AsyncSession
    from typing import Annotated
    
    async def get_db() -> AsyncGenerator[AsyncSession, None]:
        async with async_session() as session:
            yield session
    
    DB = Annotated[AsyncSession, Depends(get_db)]
    
    @app.get("/users/{user_id}")
    async def get_user(user_id: int, db: DB):
        result = await db.execute(select(User).where(User.id == user_id))
        user = result.scalar_one_or_none()
        if not user:
            raise HTTPException(status_code=404)
        return user
    ```
    
    ## Quick Reference
    
    | Operation | SQLAlchemy 2.0 Style |
    |-----------|---------------------|
    | Select all | `select(User)` |
    | Filter | `.where(User.id == 1)` |
    | First | `.scalar_one_or_none()` |
    | All | `.scalars().all()` |
    | Count | `select(func.count(User.id))` |
    | Join | `.join(Post)` |
    | Eager load | `.options(selectinload(User.posts))` |
    
    ## Additional Resources
    
    - `./references/sqlalchemy-async.md` - Async patterns, session management
    - `./references/connection-pooling.md` - Pool configuration, health checks
    - `./references/transactions.md` - Transaction patterns, isolation levels
    - `./references/migrations.md` - Alembic setup, migration strategies
    
    ## Assets
    
    - `./assets/alembic.ini.template` - Alembic configuration template
    
    ---
    
    ## See Also
    
    **Prerequisites:**
    - `python-typing-ops` - Mapped types and annotations
    - `python-async-ops` - Async database sessions
    
    **Related Skills:**
    - `python-fastapi-ops` - Dependency injection for DB sessions
    - `python-pytest-ops` - Database fixtures and testing
    

Comments (0)

Sign in to join the conversation.

No comments yet.

Reviews (0)

No reviews yet.

Related