| name | cqrs-pattern |
| description | Implement Command Query Responsibility Segregation separating read and write models. Outputs command handlers, query handlers, read model synchronization, and eventual consistency patterns. |
| argument-hint | ["domain","read/write ratio","consistency requirements","persistence technology"] |
| allowed-tools | Read, Write, Bash |
CQRS Pattern
CQRS separates write operations (commands) from read operations (queries) into distinct models. The write model enforces business rules; the read model is optimized for queries. This enables independent scaling, tailored data models for each side, and better performance for read-heavy systems.
When to Use CQRS
Good fit:
- Read/write ratio >10:1 (optimize reads independently)
- Complex domain logic on writes, simple projections on reads
- Different consistency requirements (strong writes, eventual reads)
- Event sourcing (CQRS pairs naturally with it)
Not worth it:
- Simple CRUD applications
- Small teams — operational complexity outweighs benefits
- Strong consistency required everywhere
Output Format
Command Side
from dataclasses import dataclass
from typing import Optional
@dataclass(frozen=True)
class PlaceOrderCommand:
user_id: str
items: tuple
shipping_address: dict
idempotency_key: str
@dataclass(frozen=True)
class CancelOrderCommand:
order_id: str
user_id: str
reason: str
@dataclass(frozen=True)
class UpdateShippingAddressCommand:
order_id: str
user_id: str
new_address: dict
import uuid
from dataclasses import dataclass
@dataclass
class CommandResult:
success: bool
aggregate_id: str
error: Optional[str] = None
class PlaceOrderCommandHandler:
def __init__(self, order_repo, inventory_service, event_bus):
self.order_repo = order_repo
self.inventory_service = inventory_service
self.event_bus = event_bus
async def handle(self, cmd: PlaceOrderCommand) -> CommandResult:
existing = await self.order_repo.find_by_idempotency_key(cmd.idempotency_key)
if existing:
return CommandResult(success=True, aggregate_id=existing.id)
for item in cmd.items:
if not await self.inventory_service.is_available(item["product_id"], item["quantity"]):
return CommandResult(success=False, aggregate_id="", error="Item out of stock")
order_id = (uuid.uuid4())
order = OrderAggregate(order_id)
order.place(
user_id=cmd.user_id,
items=(cmd.items),
total_cents=(i[] * i[] i cmd.items),
shipping_address=cmd.shipping_address,
)
.order_repo.save(order, idempotency_key=cmd.idempotency_key)
event order.uncommitted_events():
.event_bus.publish(event)
CommandResult(success=, aggregate_id=order_id)
:
():
.order_repo = order_repo
.event_bus = event_bus
() -> CommandResult:
order = .order_repo.load(cmd.order_id)
order:
CommandResult(success=, aggregate_id=, error=)
order.user_id != cmd.user_id:
CommandResult(success=, aggregate_id=, error=)
:
order.cancel(reason=cmd.reason, cancelled_by=cmd.user_id)
InvalidStateTransition e:
CommandResult(success=, aggregate_id=cmd.order_id, error=(e))
.order_repo.save(order)
event order.uncommitted_events():
.event_bus.publish(event)
CommandResult(success=, aggregate_id=cmd.order_id)
Query Side (Read Models)
from dataclasses import dataclass
from typing import Optional
@dataclass(frozen=True)
class GetOrderQuery:
order_id: str
user_id: str
@dataclass(frozen=True)
class ListUserOrdersQuery:
user_id: str
status: Optional[str] = None
page: int = 1
page_size: int = 20
@dataclass(frozen=True)
class GetOrderSummaryQuery:
user_id: str
days: int = 30
class GetOrderQueryHandler:
def __init__(self, read_db):
self.db = read_db
async def handle(self, query: GetOrderQuery) -> Optional[dict]:
"""Read from denormalized read model — no joins needed."""
row = await self.db.fetchrow(
"""
SELECT
o.order_id, o.status, o.total_cents, o.created_at,
o.user_id, o.payment_id, o.tracking_number,
o.items, -- Pre-serialized JSON
u.display_name as user_name, -- Pre-joined at write time
u.email as user_email
FROM order_read_model o
JOIN user_snapshot u ON u.user_id = o.user_id
WHERE o.order_id = $1
""",
query.order_id
)
if not row:
return None
if row["user_id"] != query.user_id:
return None
return dict(row)
class ListUserOrdersQueryHandler:
def __init__(self, read_db):
self.db = read_db
async def handle(self, query: ListUserOrdersQuery) -> :
conditions = []
params = [query.user_id]
query.status:
conditions.append()
params.append(query.status)
where = .join(conditions)
offset = (query.page - ) * query.page_size
rows = .db.fetch(
,
*params
)
total = .db.fetchval(
,
*params
)
{
: [(r) r rows],
: total,
: query.page,
: query.page_size,
}
:
():
.db = analytics_db
() -> :
.db.fetchrow(
,
query.user_id, query.days
)
Read Model Synchronization
import asyncio
class OrderReadModelUpdater:
"""
Subscribes to domain events and updates read models.
Must be idempotent — may receive duplicate events.
"""
def __init__(self, read_db, event_bus):
self.db = read_db
self.event_bus = event_bus
async def start(self):
await self.event_bus.subscribe(
topics=["OrderPlaced", "PaymentProcessed", "OrderShipped", "OrderCancelled"],
handler=self.handle_event
)
async def handle_event(self, event_type: str, event_data: dict, event_id: str):
already_processed = await self.db.fetchval(
"SELECT 1 FROM processed_events WHERE event_id = $1",
event_id
)
if already_processed:
return
async with self.db.transaction():
if event_type == "OrderPlaced":
await self.db.execute(
"""
INSERT INTO order_read_model (
order_id, user_id, status, total_cents,
items, items_count, created_at, updated_at
) VALUES ($1, $2, 'pending', $3, $4, $5, $6, $6)
ON CONFLICT (order_id) DO NOTHING
""",
event_data[],
event_data[],
event_data[],
json.dumps(event_data[]),
(event_data[]),
event_data[],
)
event_type == :
.db.execute(
,
event_data[],
event_data[],
event_data[],
)
event_type == :
.db.execute(
,
event_data[],
event_data[],
event_data[],
event_data[],
)
event_type == :
.db.execute(
,
event_data[],
event_data[],
event_data[],
)
.db.execute(
,
event_id
)
API Layer (Command/Query Dispatch)
from fastapi import APIRouter, Depends, HTTPException
router = APIRouter(prefix="/orders")
@router.post("/", status_code=201)
async def create_order(
request: CreateOrderRequest,
user: User = Depends(get_current_user),
handler: PlaceOrderCommandHandler = Depends(get_place_order_handler),
):
result = await handler.handle(PlaceOrderCommand(
user_id=user.id,
items=tuple(request.items),
shipping_address=request.shipping_address,
idempotency_key=request.idempotency_key or str(uuid.uuid4()),
))
if not result.success:
raise HTTPException(status_code=400, detail=result.error)
return {"order_id": result.aggregate_id}
@router.get("/{order_id}")
async def get_order(
order_id: str,
user: User = Depends(get_current_user),
handler: GetOrderQueryHandler = Depends(get_order_query_handler),
):
result = await handler.handle(GetOrderQuery(order_id=order_id, user_id=user.id))
if not result:
raise HTTPException(status_code=404)
return result
@router.get()
():
handler.handle(ListUserOrdersQuery(
user_id=user.,
status=status,
page=page,
))
Read Model Schema
CREATE TABLE order_read_model (
order_id UUID PRIMARY KEY,
user_id VARCHAR(255) NOT NULL,
status VARCHAR(50) NOT NULL,
total_cents INTEGER NOT NULL,
items JSONB NOT NULL,
items_count INTEGER NOT NULL,
payment_id VARCHAR(255),
tracking_number VARCHAR(255),
carrier VARCHAR(100),
cancel_reason TEXT,
created_at TIMESTAMPTZ NOT NULL,
updated_at TIMESTAMPTZ NOT NULL
);
CREATE INDEX idx_orders_rm_user_status ON order_read_model(user_id, status);
CREATE INDEX idx_orders_rm_created ON order_read_model(created_at DESC);
CREATE INDEX idx_orders_rm_status ON order_read_model(status) WHERE status = 'pending';
CREATE TABLE processed_events (
event_id VARCHAR(255) PRIMARY KEY,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
Rules
- Commands are intentions, queries are questions — commands may fail; queries always return something.
- Read models are denormalized — don't normalize read-side; optimize for query patterns, not storage.
- Eventual consistency on the read side — accept that reads may lag writes by milliseconds to seconds.
- Idempotent read model updaters — events may be delivered more than once; updates must be safe.
- Don't share the database — write model and read model can use different databases (write: Postgres, read: Redis or Elasticsearch).
- Commands return IDs, not entities — the command result is an ID; clients query the read model for current state.
- Multiple read models from one event stream — build as many read models as query patterns require.
- Rebuild read models when needed — they're disposable; replay events to rebuild when projection logic changes.
- CQRS ≠ event sourcing — they pair well, but either can exist without the other.
- Don't add CQRS prematurely — start simple, extract when read/write patterns clearly diverge.