
はじめに
イベント ソーシングと CQRS は、多くの場合連携して行われる 2 つのパターンであり、マイクロサービスにおけるデータの一貫性、監査証跡、パフォーマンスの最適化といった複雑な問題の解決に役立ちます。
1. イベントソーシング
1.1 コンセプト
現在の状態を保存する代わりに、発生したイベントのシーケンス全体を保存します。
Traditional (State-based):
┌──────────────────────────────┐
│ Orders Table │
│ id: O-001 │
│ status: shipped ← Chỉ biết state hiện tại
│ total: 500,000 │
│ updated_at: 2026-03-31 │
└──────────────────────────────┘
Event Sourcing:
┌──────────────────────────────────────────────────────────┐
│ Event Store (append-only) │
├────┬──────────────────┬──────────────┬──────────────────┤
│ # │ Event Type │ Data │ Timestamp │
├────┼──────────────────┼──────────────┼──────────────────┤
│ 1 │ OrderCreated │ {items, ...} │ 10:00:00 │
│ 2 │ PaymentReceived │ {amount} │ 10:01:00 │
│ 3 │ ItemsReserved │ {items} │ 10:01:05 │
│ 4 │ OrderShipped │ {tracking} │ 10:30:00 │
└────┴──────────────────┴──────────────┴──────────────────┘
Current State = replay(events) → Order{status: "shipped"}
1.2 イベントストア
Đặc điểm:
├── Append-only: Không bao giờ update hoặc delete events
├── Immutable: Events là facts đã xảy ra, không thể thay đổi
├── Ordered: Events có thứ tự rõ ràng (sequence number)
└── Stream: Events được nhóm theo aggregate (ví dụ: order-O-001)
Implementation options:
├── EventStoreDB (purpose-built, recommended)
├── PostgreSQL + events table
├── Apache Kafka (log-based)
└── DynamoDB Streams (AWS)
PostgreSQL イベント ストア:
CREATE TABLE events (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
stream_id VARCHAR(255) NOT NULL, -- "order-O-001"
version BIGINT NOT NULL, -- sequence number
event_type VARCHAR(255) NOT NULL, -- "OrderCreated"
data JSONB NOT NULL, -- event payload
metadata JSONB DEFAULT '{}', -- traceId, userId, ...
created_at TIMESTAMPTZ DEFAULT NOW(),
UNIQUE(stream_id, version) -- đảm bảo ordering
);
CREATE INDEX idx_events_stream ON events(stream_id, version);
1.3 状態の再構築
def get_order(order_id: str) -> Order:
events = event_store.get_events(stream_id=f"order-{order_id}")
order = Order() # empty state
for event in events:
order.apply(event) # replay từng event
return order # current state
class Order:
def apply(self, event):
match event.type:
case "OrderCreated":
self.id = event.data["id"]
self.status = "created"
self.items = event.data["items"]
case "PaymentReceived":
self.status = "paid"
case "OrderShipped":
self.status = "shipped"
self.tracking = event.data["tracking"]
1.4 スナップショットの最適化
ストリームにイベントが多すぎて(数千)、再生が遅い場合 → スナップショットを使用します。
Event Stream cho order-O-001:
Event 1: OrderCreated
Event 2: ItemAdded
...
Event 500: ItemRemoved
──── Snapshot at version 500 ────
{ status: "processing", items: [...], total: 1000000 }
Event 501: PaymentReceived
Event 502: OrderShipped
Rebuild: Load snapshot (v500) + replay events 501-502
→ Nhanh hơn nhiều so với replay 502 events
1.5 メリットとデメリット
✅ Ưu điểm:
├── Complete audit trail (ai làm gì, lúc nào)
├── Time travel: Rebuild state tại bất kỳ thời điểm
├── Event replay: Fix bug rồi replay events để sửa data
├── Natural fit cho event-driven architecture
└── Debug: Hiểu chính xác điều gì đã xảy ra
❌ Nhược điểm:
├── Complexity: Khó hơn CRUD đáng kể
├── Query: Không thể query trực tiếp (cần CQRS)
├── Schema evolution: Thay đổi event format phức tạp
├── Storage: Nhiều events → nhiều storage
└── Learning curve: Team cần thời gian adapt
2. CQRS — コマンドクエリの責任の分離
2.1 概念
書き込み用モデル (コマンド) と 読み取り用モデル (クエリ) を分けます。
Traditional:
Client ──CRUD──▶ Same Model ──▶ Same Database
CQRS:
┌─────────────────────────────────────┐
│ API Layer │
└──────────┬───────────────┬───────────┘
│ │
┌──────────▼──────┐ ┌──────▼──────────┐
│ Command Side │ │ Query Side │
│ (Write Model) │ │ (Read Model) │
│ │ │ │
│ - CreateOrder │ │ - GetOrderDetails│
│ - CancelOrder │ │ - ListOrders │
│ - UpdateStatus │ │ - SearchOrders │
└────────┬────────┘ └────────▲─────────┘
│ │
┌────────▼────────┐ ┌───────┴─────────┐
│ PostgreSQL │ │ Elasticsearch │
│ (Write DB) │ │ (Read DB) │
│ Normalized │ │ Denormalized │
└────────┬────────┘ └─────────────────┘
│ ▲
└───── Events ──────┘
(sync read model)
2.2 なぜ読み取りと書き込みを分けるのでしょうか?
Write: Read:
├── Ít operations hơn ├── Nhiều operations hơn (10:1 ratio)
├── Cần ACID consistency ├── Eventual consistency OK
├── Normalized schema ├── Denormalized, pre-joined
├── Complex validation ├── Simple query, fast response
├── Scale: moderate ├── Scale: aggressive (caching, replicas)
└── PostgreSQL (optimal) └── Elasticsearch/Redis (optimal)
2.3 読み取りモデルの同期
Option 1: Domain Events (khuyến nghị)
Write DB ──event──▶ Kafka ──▶ Read Model Updater ──▶ Read DB
Option 2: Change Data Capture (CDC)
Write DB ──Debezium──▶ Kafka ──▶ Read Model Updater ──▶ Read DB
Option 3: Dual Write (KHÔNG khuyến nghị)
Service ──write──▶ Write DB
──write──▶ Read DB ← Có thể inconsistent!
2.4 CQRS + イベントソーシング
両方のパターンを組み合わせます。
Command Flow:
Client ──CreateOrder──▶ Command Handler
│
Validate
│
Append event to Event Store
│
Publish event to Kafka
│
▼
Event Store (source of truth)
Query Flow:
Kafka ──OrderCreated──▶ Projection Handler
│
Update Read Model (Elasticsearch)
│
Client ──GetOrder──▶ Query Handler ──▶ Read from Elasticsearch
2.5 最終的な整合性
Timeline:
T0: Client tạo order (write to Event Store)
T1: Event published to Kafka (~5ms)
T2: Projection handler updates Elasticsearch (~50ms)
T3: Read model available (~100ms after T0)
Giữa T0 và T3: Read model chưa có data mới = Eventual Consistency
Giải pháp UX:
├── Optimistic UI: Client hiển thị ngay sau write, không đợi read model
├── Read-your-writes: Sau write, query write DB cho user đó
├── Polling/WebSocket: Client poll cho đến khi read model updated
└── Inbox pattern: Return 202 Accepted + polling endpoint
3. イベント ソーシング/CQRS をいつ使用するか?
3.1 次の場合にはイベント ソーシングを使用する必要があります。
- ✅ 完全な 監査証跡 が必要 (財務、医療、法律)
- ✅ タイムトラベルが必要 (いつでも状態を再構築)
- ✅ ドメインには 複雑な状態遷移があります (注文ワークフロー、予約)
- ✅ 本番環境の問題をデバッグする必要があります (イベントをリプレイ)
- ✅ イベント駆動型アーキテクチャはすでに基盤となっています
3.2 CQRS はどのような場合に使用する必要がありますか?
- ✅ 読み取り/書き込み比 大きな違い (10:1 以上)
- ✅ 読み取りと書き込みには 異なるスケールが必要です
- ✅ 非正規化/事前計算が必要なモデルを読み取る
- ✅ 複雑なクエリには 検索エンジン (Elasticsearch) が必要です
3.3 次の場合には使用しないでください。
- ❌ シンプルな CRUD アプリケーション
- ❌ チームにはイベントドリブンの経験がない
- ❌ ドメインが単純で、状態遷移が少ない
- ❌ 一貫性要件 = どこでも強い一貫性
- ❌ 締め切りが迫っており、すぐに発送する必要があります
4. まとめ
| パターン | キーポイント |
|---|---|
| イベントソーシング | 状態の代わりにイベントを保存、追加のみ、監査証跡 |
| イベントストア | 不変のイベントログ、信頼できる情報源 |
| スナップショット | 定期的なスナップショットを使用して再構築を最適化する |
| CQRS | 個別の読み取り/書き込みモデル、個別にスケーリング |
| 最終的な整合性 | 書き込み後のモデル更新の読み取り (ms レベルの遅延) |
| 投影 | イベントの処理 → 読み取りモデルの更新 |
次の記事: Saga パターン — 各サービスが独自のデータベースを持つ場合の分散トランザクションの処理。