LabHub

ブログ

CQRSパターン実践実装ガイド

한국어English日本語

CQRS Pattern Implementation

1. CQRSとは?

CQRS(Command Query Responsibility Segregation)は、データの書き込み(Command)と読み取り(Query)を別々のモデルに分離するアーキテクチャパターンです。

従来のCRUD vs CQRS

従来のCRUDでは、1つのモデルが読み取りと書き込みの両方を担当します。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の複雑さがメリットを上回る場合 チームの経験不足: Event SourcingやメッセージキューOの経験なしに導入するのはリスクが高い

クイズ

Q1: 「CQRSパターン実践実装ガイド」の主なトピックは何ですか? Command/Query分離の原理からEvent Sourcing連携、Kafka活用、Spring Bootベースの実践実装例まで、CQRSパターンを完全網羅します。

Q2: CQRSとは?とは何ですか? CQRS(Command Query Responsibility Segregation)は、データの書き込み(Command)と読み取り(Query)を別々のモデルに分離するアーキテクチャパターンです。 従来のCRUD vs CQRS 従来のCRUDでは、1つのモデルが読み取りと書き込みの両方を担当します。CQRSではこれを分離し、それぞれを最適化します。 CQRSを適用すべきタイミングは?

Q3: Commandモデルの実装の核心的な概念を説明してください。 Commandの定義 Command Handler

Q4: Queryモデルの実装の主な特徴は何ですか? Queryモデルは、読み取りに最適化された別のデータストアを使用します。

コメント

まだコメントはありません。

ログインするとコメントできます