Chuyển đến nội dung chính

レッスン 9: イベント ソーシングと CQRS

イベント ソーシング - 状態をイベント文字列として保存、イベント ストア、スナップショットの最適化。 CQRS — 個別のコマンド モデルとクエリ モデル、結果整合性、個別の読み取り/書き込みデータベース、CQRS をいつ使用するか。

🏗️ アーキテクチャ — レッスン 9 レッスン 9: イベント ソーシングと CQRS

クラウドネイティブのマイクロサービスアーキテクチャ

パート 3: マイクロサービスにおけるデータ管理

xdev.asia

レッスン 9: イベント ソーシングと CQRS

はじめに

イベント ソーシングと 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 パターン — 各サービスが独自のデータベースを持つ場合の分散トランザクションの処理。