LabHub

블로그

CQRS 패턴 실전 구현 가이드

한국어English日本語

CQRS Pattern Implementation

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를 적용하나?

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의 복잡성이 이점보다 클 때 팀 경험 부족: 이벤트 소싱, 메시지 큐 경험 없이 도입하면 위험

댓글

아직 댓글이 없습니다.

로그인하면 댓글을 쓸 수 있습니다