MongoDB with Motor
Motor is the official async MongoDB driver for Python. It uses the same API as pymongo but with async/await.
Setup
pip install motor
from motor.motor_asyncio import AsyncIOMotorClient
MONGO_URL = "mongodb://localhost:27017"
client: AsyncIOMotorClient | None = None
async def connect_db():
global client
client = AsyncIOMotorClient(MONGO_URL)
async def close_db():
if client:
client.close()
def get_db():
return client["techblog"]
Lifespan Integration
from contextlib import asynccontextmanager
from fastapi import FastAPI
from app.database import connect_db, close_db
@asynccontextmanager
async def lifespan(app: FastAPI):
await connect_db()
yield
await close_db()
app = FastAPI(lifespan=lifespan)
CRUD Operations
from fastapi import APIRouter, HTTPException, status
from app.database import get_db
from bson import ObjectId
from pydantic import BaseModel, Field
from typing import Optional
router = APIRouter(prefix="/items", tags=["items"])
class ItemCreate(BaseModel):
name: str
price: float
description: Optional[str] = None
class ItemResponse(BaseModel):
id: str = Field(alias="_id")
name: str
price: float
description: Optional[str] = None
@router.post("/", status_code=status.HTTP_201_CREATED)
async def create_item(payload: ItemCreate):
db = get_db()
result = await db.items.insert_one(payload.model_dump())
created = await db.items.find_one({"_id": result.inserted_id})
created["_id"] = str(created["_id"])
return created
@router.get("/{item_id}")
async def get_item(item_id: str):
db = get_db()
try:
obj_id = ObjectId(item_id)
except Exception:
raise HTTPException(status_code=400, detail="Invalid ID format")
item = await db.items.find_one({"_id": obj_id})
if not item:
raise HTTPException(status_code=404, detail="Item not found")
item["_id"] = str(item["_id"])
return item
@router.get("/")
async def list_items(skip: int = 0, limit: int = 20):
db = get_db()
items = []
cursor = db.items.find().skip(skip).limit(limit)
async for doc in cursor:
doc["_id"] = str(doc["_id"])
items.append(doc)
return items
@router.put("/{item_id}")
async def update_item(item_id: str, payload: ItemCreate):
db = get_db()
try:
obj_id = ObjectId(item_id)
except Exception:
raise HTTPException(status_code=400, detail="Invalid ID format")
result = await db.items.replace_one({"_id": obj_id}, payload.model_dump())
if result.matched_count == 0:
raise HTTPException(status_code=404, detail="Item not found")
updated = await db.items.find_one({"_id": obj_id})
updated["_id"] = str(updated["_id"])
return updated
@router.delete("/{item_id}", status_code=status.HTTP_204_NO_CONTENT)
async def delete_item(item_id: str):
db = get_db()
try:
obj_id = ObjectId(item_id)
except Exception:
raise HTTPException(status_code=400, detail="Invalid ID format")
result = await db.items.delete_one({"_id": obj_id})
if result.deleted_count == 0:
raise HTTPException(status_code=404, detail="Item not found")
MongoDB Schema Notes
| Aspect | MongoDB Behavior |
|---|
| ID field | _id is an ObjectId - convert to string for JSON |
| Schema | Schema-less - documents in the same collection can have different fields |
| Indexes | Create with db.collection.create_index("field") |
| Relations | Use $lookup aggregation or manual references (no JOINs) |
| Transactions | Supported since MongoDB 4.0 (requires replica set) |
Redis with redis-py
redis-py is the official Python client. Use aioredis (merged into redis-py 4.x+) for async support.
Setup
pip install redis
from redis.asyncio import Redis
redis_client: Redis | None = None
async def init_redis():
global redis_client
redis_client = Redis(host="localhost", port=6379, decode_responses=True)
async def close_redis():
if redis_client:
await redis_client.aclose()
def get_redis() -> Redis:
return redis_client
Lifespan Integration
@asynccontextmanager
async def lifespan(app: FastAPI):
await connect_db()
await init_redis()
yield
await close_redis()
await close_db()
Redis CRUD
router = APIRouter(prefix="/sessions", tags=["sessions"])
@router.post("/{session_id}")
async def set_session(session_id: str, data: dict):
r = get_redis()
await r.hset(f"session:{session_id}", mapping=data)
await r.expire(f"session:{session_id}", 3600)
return {"status": "ok"}
@router.get("/{session_id}")
async def get_session(session_id: str):
r = get_redis()
data = await r.hgetall(f"session:{session_id}")
if not data:
raise HTTPException(status_code=404, detail="Session not found")
return data
@router.delete("/{session_id}", status_code=204)
async def delete_session(session_id: str):
r = get_redis()
await r.delete(f"session:{session_id}")
Caching Patterns
Pattern 1: Cache-Aside
@router.get("/items/{item_id}")
async def get_item_with_cache(item_id: str):
r = get_redis()
cache_key = f"item:{item_id}"
cached = await r.get(cache_key)
if cached:
return {"source": "cache", "data": json.loads(cached)}
db = get_db()
try:
obj_id = ObjectId(item_id)
except Exception:
raise HTTPException(status_code=400, detail="Invalid ID")
item = await db.items.find_one({"_id": obj_id})
if not item:
raise HTTPException(status_code=404, detail="Item not found")
item["_id"] = str(item["_id"])
await r.setex(cache_key, 300, json.dumps(item, default=str))
return {"source": "database", "data": item}
Pattern 2: Invalidation on Write
@router.put("/items/{item_id}")
async def update_item_clear_cache(item_id: str, payload: ItemCreate):
db = get_db()
r = get_redis()
try:
obj_id = ObjectId(item_id)
except Exception:
raise HTTPException(status_code=400, detail="Invalid ID")
result = await db.items.replace_one({"_id": obj_id}, payload.model_dump())
if result.matched_count == 0:
raise HTTPException(status_code=404, detail="Item not found")
await r.delete(f"item:{item_id}")
updated = await db.items.find_one({"_id": obj_id})
updated["_id"] = str(updated["_id"])
return updated
Pattern 3: Rate Limiting
from fastapi import Request
@app.middleware("http")
async def rate_limit_middleware(request: Request, call_next):
r = get_redis()
client_ip = request.client.host
key = f"ratelimit:{client_ip}"
count = await r.incr(key)
if count == 1:
await r.expire(key, 60)
if count > 100:
return JSONResponse(status_code=429, content={"detail": "Too many requests"})
return await call_next(request)
Redis Command Reference
| Command | Python | Use Case |
|---|
SET key value [EX seconds] | await r.set(key, val, ex=ttl) | String cache |
GET key | await r.get(key) | Read cache |
HSET key field value | await r.hset(key, mapping={}) | Object store |
HGETALL key | await r.hgetall(key) | Read object |
EXPIRE key seconds | await r.expire(key, ttl) | Set TTL |
INCR key | await r.incr(key) | Counters |
DEL key | await r.delete(key) | Remove key |
PUBLISH channel msg | await r.publish(ch, msg) | Pub/Sub |
SQL vs NoSQL - How to Choose
| Factor | SQL (PostgreSQL) | NoSQL (MongoDB) | Redis |
|---|
| Schema | Rigid, enforced | Flexible, no enforcement | Key-value / struct |
| Relations | JOINs, foreign keys | Embedded docs, $lookup | Manual |
| ACID | Full ACID | Multi-doc ACID (4.0+) | No (single key atomic) |
| Query | Complex joins, aggregations | Rich query language | By key only |
| Scaling | Vertical (read replicas) | Horizontal (sharding) | Horizontal (cluster) |
| Best for | Relational data, reporting | Rapid prototyping, varied docs | Caching, sessions, queues |
| FastAPI fit | SQLAlchemy + Alembic | Motor + FastAPI | redis-py + Depends |
Key Takeaways
- Use Motor (
motor.motor_asyncio) for async MongoDB - it mirrors pymongo with await on every call
- Convert
ObjectId to str before returning JSON responses; use bson.objectid.ObjectId for queries
- Use redis-py (
Redis.asyncio) for Redis - hset/hgetall for objects, setex for TTL caches
- Implement cache-aside pattern: check Redis first, fall back to DB, write to Redis on miss
- Invalidate cache on writes to prevent stale data; use
EXPIRE/setex for automatic eviction
- Choose SQL for complex queries and relationships; MongoDB for flexible schemas; Redis for caching and ephemeral data