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.
Virus-scanned
Reviewed automatically before listing.
Download
0xdarkmatter-claude-mods-skills_python-database-ops-3dfaf0b.zip · 12 KB
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 annotationspython-async-ops- Async database sessions
Related Skills:
python-fastapi-ops- Dependency injection for DB sessionspython-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.
Reviews (0)
No reviews yet.
No comments yet.