
1. CQRS란?
CQRS(Command Query Responsibility Segregation)는 데이터의 쓰기(Command)와 읽기(Query)를 별도의 모델로 분리하는 아키텍처 패턴입니다.
전통적 CRUD vs CQRS
전통적 CRUD에서는 하나의 모델이 읽기와 쓰기를 모두 담당합니다. CQRS에서는 이를 분리하여 각각 최적화합니다.
전통적 CRUD:
Client → [같은 모델] → Database
CQRS:
Client (쓰기) → [Command Model] → Write DB
Client (읽기) → [Query Model] → Read DB
언제 CQRS를 적용하나?
- 읽기와 쓰기의 비율이 크게 다를 때 (읽기 90%+)
- 읽기/쓰기에 서로 다른 최적화가 필요할 때
- 복잡한 도메인 로직이 있을 때
- 이벤트 소싱과 함께 사용할 때
2. Command 모델 구현
Command 정의
from dataclasses import dataclass
from datetime import datetime
from uuid import UUID, uuid4
# Command 객체
@dataclass(frozen=True)
class CreateOrderCommand:
customer_id: UUID
items: list[dict]
shipping_address: str
@dataclass(frozen=True)
class CancelOrderCommand:
order_id: UUID
reason: str
# Domain Event
@dataclass(frozen=True)
class OrderCreatedEvent:
event_id: UUID
order_id: UUID
customer_id: UUID
items: list[dict]
total_amount: float
created_at: datetime
Command Handler
from typing import Protocol
class EventStore(Protocol):
def append(self, stream_id: str, events: list) -> None: ...
def load(self, stream_id: str) -> list: ...
class OrderCommandHandler:
def __init__(self, event_store: EventStore, event_bus):
self.event_store = event_store
self.event_bus = event_bus
def handle_create_order(self, cmd: CreateOrderCommand):
order_id = uuid4()
total = sum(item['price'] * item['quantity']
for item in cmd.items)
# 비즈니스 규칙 검증
if total <= 0:
raise ValueError("Order total must be positive")
if not cmd.items:
raise ValueError("Order must have at least one item")
# 이벤트 생성
event = OrderCreatedEvent(
event_id=uuid4(),
order_id=order_id,
customer_id=cmd.customer_id,
items=cmd.items,
total_amount=total,
created_at=datetime.utcnow(),
)
# 이벤트 저장 + 발행
self.event_store.append(f"order-{order_id}", [event])
self.event_bus.publish("order.created", event)
return order_id
def handle_cancel_order(self, cmd: CancelOrderCommand):
# 기존 이벤트에서 상태 복원
events = self.event_store.load(f"order-{cmd.order_id}")
order = Order.from_events(events)
if order.status == "cancelled":
raise ValueError("Order already cancelled")
cancel_event = OrderCancelledEvent(
event_id=uuid4(),
order_id=cmd.order_id,
reason=cmd.reason,
cancelled_at=datetime.utcnow(),
)
self.event_store.append(
f"order-{cmd.order_id}", [cancel_event]
)
self.event_bus.publish("order.cancelled", cancel_event)
3. Query 모델 구현
Query 모델은 읽기에 최적화된 별도의 데이터 저장소를 사용합니다.
from dataclasses import dataclass
# 읽기 전용 모델 (Materialized View)
@dataclass
class OrderSummary:
order_id: str
customer_name: str
item_count: int
total_amount: float
status: str
created_at: str
class OrderQueryHandler:
def __init__(self, read_db):
self.read_db = read_db
def get_order(self, order_id: str) -> OrderSummary:
row = self.read_db.execute(
"SELECT * FROM order_summaries WHERE order_id = %s",
(order_id,)
)
return OrderSummary(**row)
def list_orders_by_customer(
self, customer_id: str, limit: int = 20
) -> list[OrderSummary]:
rows = self.read_db.execute(
"""SELECT * FROM order_summaries
WHERE customer_id = %s
ORDER BY created_at DESC
LIMIT %s""",
(customer_id, limit)
)
return [OrderSummary(**r) for r in rows]
def search_orders(
self, status: str = None, min_amount: float = None
) -> list[OrderSummary]:
query = "SELECT * FROM order_summaries WHERE 1=1"
params = []
if status:
query += " AND status = %s"
params.append(status)
if min_amount:
query += " AND total_amount >= %s"
params.append(min_amount)
return self.read_db.execute(query, params)
4. Event Sourcing 연동
이벤트 스토어 구현
import json
from datetime import datetime
class PostgresEventStore:
def __init__(self, conn):
self.conn = conn
def create_table(self):
self.conn.execute("""
CREATE TABLE IF NOT EXISTS event_store (
id BIGSERIAL PRIMARY KEY,
stream_id VARCHAR(255) NOT NULL,
event_type VARCHAR(255) NOT NULL,
event_data JSONB NOT NULL,
metadata JSONB DEFAULT '{}',
version INTEGER NOT NULL,
created_at TIMESTAMP DEFAULT NOW(),
UNIQUE(stream_id, version)
);
CREATE INDEX IF NOT EXISTS idx_event_stream
ON event_store(stream_id, version);
""")
def append(self, stream_id: str, events: list):
current = self._get_latest_version(stream_id)
for i, event in enumerate(events):
version = current + i + 1
self.conn.execute(
"""INSERT INTO event_store
(stream_id, event_type, event_data, version)
VALUES (%s, %s, %s, %s)""",
(stream_id, type(event).__name__,
json.dumps(event.__dict__, default=str),
version)
)
def load(self, stream_id: str) -> list:
rows = self.conn.execute(
"""SELECT event_type, event_data FROM event_store
WHERE stream_id = %s ORDER BY version""",
(stream_id,)
)
return [self._deserialize(r) for r in rows]
def _get_latest_version(self, stream_id: str) -> int:
result = self.conn.execute(
"SELECT MAX(version) FROM event_store WHERE stream_id = %s",
(stream_id,)
)
return result or 0
5. Kafka를 활용한 이벤트 동기화
from confluent_kafka import Producer, Consumer
import json
class KafkaEventBus:
def __init__(self, bootstrap_servers: str):
self.producer = Producer({
'bootstrap.servers': bootstrap_servers,
'acks': 'all',
'enable.idempotence': True,
})
def publish(self, topic: str, event):
self.producer.produce(
topic=topic,
key=str(event.order_id).encode(),
value=json.dumps(
event.__dict__, default=str
).encode(),
)
self.producer.flush()
class ReadModelProjector:
"""이벤트를 소비하여 Read DB를 업데이트"""
def __init__(self, read_db, bootstrap_servers: str):
self.read_db = read_db
self.consumer = Consumer({
'bootstrap.servers': bootstrap_servers,
'group.id': 'read-model-projector',
'auto.offset.reset': 'earliest',
'enable.auto.commit': False,
})
self.consumer.subscribe([
'order.created', 'order.cancelled'
])
def run(self):
while True:
msg = self.consumer.poll(1.0)
if msg is None:
continue
event = json.loads(msg.value())
topic = msg.topic()
if topic == 'order.created':
self._project_order_created(event)
elif topic == 'order.cancelled':
self._project_order_cancelled(event)
self.consumer.commit()
def _project_order_created(self, event):
self.read_db.execute(
"""INSERT INTO order_summaries
(order_id, customer_id, item_count,
total_amount, status, created_at)
VALUES (%s, %s, %s, %s, %s, %s)""",
(event['order_id'], event['customer_id'],
len(event['items']), event['total_amount'],
'active', event['created_at'])
)
def _project_order_cancelled(self, event):
self.read_db.execute(
"""UPDATE order_summaries
SET status = 'cancelled'
WHERE order_id = %s""",
(event['order_id'],)
)
6. API 레이어
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
app = FastAPI()
class CreateOrderRequest(BaseModel):
customer_id: str
items: list[dict]
shipping_address: str
# Command 엔드포인트 (POST/PUT/DELETE)
@app.post("/orders")
async def create_order(req: CreateOrderRequest):
cmd = CreateOrderCommand(
customer_id=UUID(req.customer_id),
items=req.items,
shipping_address=req.shipping_address,
)
order_id = command_handler.handle_create_order(cmd)
return {"order_id": str(order_id), "status": "created"}
# Query 엔드포인트 (GET)
@app.get("/orders/{order_id}")
async def get_order(order_id: str):
result = query_handler.get_order(order_id)
if not result:
raise HTTPException(status_code=404)
return result
@app.get("/customers/{customer_id}/orders")
async def list_customer_orders(customer_id: str, limit: int = 20):
return query_handler.list_orders_by_customer(customer_id, limit)
7. 퀴즈
Q1: CQRS에서 Eventual Consistency가 발생하는 이유와 해결 전략은?
Command가 Write DB에 저장된 후 이벤트가 Kafka를 통해 Read DB에 반영되기까지 시간차가 있기 때문입니다. 해결 전략:
Read-your-writes: 방금 쓴 데이터는 Write DB에서 직접 읽기 Polling: 클라이언트가 Read DB가 업데이트될 때까지 재시도 Versioning: 이벤트 버전을 추적하여 최신 여부 확인 WebSocket/SSE: 업데이트 완료 시 실시간 알림
Q2: Event Sourcing에서 이벤트 스키마가 변경되면 어떻게 처리하나요?
이벤트 업캐스팅(Upcasting) 패턴을 사용합니다:
이벤트에 버전 필드를 추가합니다 (v1, v2, ...) 이전 버전 이벤트를 로드할 때 최신 스키마로 변환하는 업캐스터를 작성합니다 기존 이벤트는 절대 수정하지 않습니다 (불변 원칙)
대안으로 스냅샷을 활용하여 모든 이벤트를 재생하지 않을 수도 있습니다.
Q3: CQRS를 적용하지 말아야 하는 경우는?
단순한 CRUD 앱: 읽기/쓰기 모델이 거의 동일할 때 강한 일관성이 필수: 실시간 잔액 확인 등 eventual consistency를 허용할 수 없을 때 소규모 프로젝트: CQRS의 복잡성이 이점보다 클 때 팀 경험 부족: 이벤트 소싱, 메시지 큐 경험 없이 도입하면 위험