
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を適用すべきタイミングは?
- 読み取りと書き込みの比率が大きく異なる場合(読み取り90%以上)
- 読み取り/書き込みにそれぞれ異なる最適化が必要な場合
- 複雑なドメインロジックがある場合
- Event Sourcingと併用する場合
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モデルは、読み取りに最適化された別のデータストアを使用します。