Skip to content

CQRS

実装: messaging/cqrs/ / 実行: go test ./messaging/cqrs/

普通のアプリは1つのモデルで読み書きの両方をこなす。だが読みと書きは要求が違う。書きは正しさが命で、読みは速さと使いやすい形が命だ。両立させると両方が中途半端になる。CQRSはこれを分ける。書き込み側は状態を検証してイベントを流し、読み取り側はそのイベントを畳んで用途別のビューを作る。1つのイベント列から注文状況・顧客ごとの購入額・総売上を独立に組める。分離を実装し、代償の結果整合まで確かめる。

この章で作るもの

イベントソーシングで、状態を保存せずイベントから導く仕組みを作った。CQRS(Command Query Responsibility Segregation)は、その導き方を「書き込み」と「読み取り」で分ける考え方だ。多くのアプリは、1 つのデータモデルで読み書きの両方をこなす。同じテーブル、同じオブジェクトに書き、そこから読む。だが読みと書きは、求めるものがまるで違う。

書き込みが大事にするのは正しさだ。「支払済みの注文は取り消せない」「在庫を超える注文は受けない」といった不変条件を守り、状態を矛盾なく変える。一方、読み取りが大事にするのは速さと形だ。「この顧客の購入履歴を一覧で」「今日の売上合計を」といった問いに、素早く、都合のよい形で答えたい。1 つのモデルで両方を最適化しようとすると、書きのための正規化された構造が読みには使いにくく、読みのための非正規化が書きの整合性を危うくする。どちらも中途半端になる。CQRS は、書き込みモデルと読み取りモデルを分けてこれを解く。書き込み側はイベントを流し、読み取り側はそのイベントから、用途ごとに最適な読みモデル(射影)を組む。

コマンド ──▶ 書き込み側 ──イベント──▶ 読み取り側 ┬─▶ 注文状況の一覧
        (検証・整合性)     (追記ログ)        ├─▶ 顧客ごとの購入額
                                              └─▶ 総売上
              1 つのイベント列から、用途別のビューを何個でも
CQRS。書き込み側はコマンドを検証してイベントを流す。読み取り側は同じイベントから、用途別の読みモデルをいくつも組み立てる

順に見ていく。

  1. 書きと読みを分ける: 書き込み側は正しさ、読み取り側は速さと形。それぞれを独立に最適化できる
  2. 1 イベント列から複数ビュー: 同じイベントを別々に畳んで、用途ごとの読みモデルをいくつも作る
  3. 代償は結果整合: 読みモデルは書きに少し遅れて追いつく。その間は古い値を返す

① 書き込み側: コマンドを検証してイベントを流す

まず書き込み側だ。コマンド(注文する、支払う、取り消す)を受け、状態を検証してからイベントを流す。検証のための状態は、イベントを畳んで導く(イベントソーシングと同じ):

go

var (
	// ErrExists は既に存在する注文を再度作ろうとしたとき。
	ErrExists = errors.New("cqrs: order already exists")
	// ErrNotFound は無い注文を操作しようとしたとき。
	ErrNotFound = errors.New("cqrs: order not found")
	// ErrNotPayable は支払えない状態(既に支払済み/取消済み)。
	ErrNotPayable = errors.New("cqrs: order not payable")
	// ErrPaidCannotCancel は支払済みは取り消せない。
	ErrPaidCannotCancel = errors.New("cqrs: paid order cannot be cancelled")
)

// WriteSide はコマンドを受け、検証して、イベントを流す(書き込みモデル)。
// 状態はイベントから毎回導く(検証のためだけに使う最小の再構成)。
type WriteSide struct{ store *Store }

// NewWriteSide は書き込み側を作る。
func NewWriteSide(store *Store) *WriteSide { return &WriteSide{store: store} }

// statusOf は id の現在状態をイベントから導く(空文字なら未存在)。
func (w *WriteSide) statusOf(id string) string {
	st := ""
	for _, e := range w.store.events {
		if e.OrderID != id {
			continue
		}
		switch e.Kind {
		case Placed:
			st = "placed"
		case Paid:
			st = "paid"
		case Cancelled:
			st = "cancelled"
		}
	}
	return st
}

// Place は注文を作る。既存 ID は拒否。
func (w *WriteSide) Place(id, customer string, amount int) error {
	if w.statusOf(id) != "" {
		return ErrExists
	}
	w.store.append(Event{Kind: Placed, OrderID: id, Customer: customer, Amount: amount})
	return nil
}

// Pay は注文を支払う。placed の注文だけが支払える。
func (w *WriteSide) Pay(id string) error {
	switch w.statusOf(id) {
	case "":
		return ErrNotFound
	case "placed":
		w.store.append(Event{Kind: Paid, OrderID: id})
		return nil
	default:
		return ErrNotPayable
	}
}

// Cancel は注文を取り消す。支払済みは取り消せない。
func (w *WriteSide) Cancel(id string) error {
	switch w.statusOf(id) {
	case "":
		return ErrNotFound
	case "paid":
		return ErrPaidCannotCancel
	default:
		w.store.append(Event{Kind: Cancelled, OrderID: id})
		return nil
	}
}

書き込み側の仕事は、不正な変更を防ぐことだ。Pay は placed の注文だけを支払える。既に支払済み・取消済みなら ErrNotPayable で弾く。Cancel は支払済みを取り消せない。これらの不変条件を守るのが書き込みモデルの責務で、ここでは形の美しさより正しさが優先される。検証を通ったコマンドだけが、イベントとしてログに流れる。テストで、重複作成・二重支払い・支払済みの取消がすべて拒否されることを固定した。

② 読み取り側: 同じイベントから複数のビューを組む

次が CQRS の醍醐味だ。読み取り側は、書き込み側が流した同じイベント列を畳んで、用途ごとに別々の読みモデルを作る。注文状況の一覧、顧客ごとの購入額、総売上。1 つのイベント列から、何種類でも:

go

// ReadSide は同じイベント列から複数の読みモデル(射影)を組み立てる。
// cursor までを処理済みとし、CatchUp で新しいイベントに追いつく。
type ReadSide struct {
	store  *Store
	cursor int // どのバージョンまで処理したか

	// 用途別の読みモデル。それぞれ同じイベントを別の形に畳む。
	Status  map[string]string // 注文ID → 状態(注文状況の一覧)
	Spend   map[string]int    // 顧客 → 支払済み合計(顧客ごとの購入額)
	Revenue int               // 総売上(支払済みの合計)
	orders  map[string]order  // 射影が必要とする注文明細(内部)
}

type order struct {
	customer string
	amount   int
}

// NewReadSide は空の読みモデルで読み取り側を作る。
func NewReadSide(store *Store) *ReadSide {
	return &ReadSide{
		store:  store,
		Status: make(map[string]string),
		Spend:  make(map[string]int),
		orders: make(map[string]order),
	}
}

// CatchUp は未処理のイベントを畳んで読みモデルを最新にする。
// これを呼ぶまで読みモデルは書き込みに遅れる(結果整合)。
func (r *ReadSide) CatchUp() {
	for _, e := range r.store.events {
		if e.Version <= r.cursor {
			continue
		}
		r.apply(e)
		r.cursor = e.Version
	}
}

// Lag は読みモデルが書き込みにどれだけ遅れているか(未処理イベント数)。
func (r *ReadSide) Lag() int { return len(r.store.events) - r.cursor }

func (r *ReadSide) apply(e Event) {
	switch e.Kind {
	case Placed:
		r.Status[e.OrderID] = "placed"
		r.orders[e.OrderID] = order{customer: e.Customer, amount: e.Amount}
	case Paid:
		r.Status[e.OrderID] = "paid"
		o := r.orders[e.OrderID]
		r.Spend[o.customer] += o.amount // 顧客ごとの購入額に反映
		r.Revenue += o.amount           // 総売上に反映
	case Cancelled:
		r.Status[e.OrderID] = "cancelled"
	}
}

apply を見ると、1 つのイベントが 3 つの読みモデルを同時に更新している。Paid イベントは、注文状況を paid にし、その顧客の購入額を増やし、総売上に加える。同じ出来事を、それぞれのビューが必要な形で受け取る。ここが 1 つのモデルでは難しいところだ。「顧客ごとの購入額」に最適な形と「総売上」に最適な形と「注文状況一覧」に最適な形は、それぞれ違う。CQRS ではビューごとに専用の読みモデルを持てるので、どれも読みやすい形にできる。新しい問い(「日別の売上は」)が増えても、同じイベントを別の畳み方で射影するだけでいい。書き込み側は一切変えない。テストで、同じイベント列から 3 つの読みモデルが正しく構築されることを固定した。

③ 結果整合: 読みは書きに遅れる

分離には代償がある。読みモデルは、書き込みと同じ瞬間には更新されない。書き込み側がイベントを流し、読み取り側がそれを処理して初めて、読みモデルに反映される。この処理が追いつくまで、読みモデルは少し古い値を返す。それを担うのが、上のコードにある CatchUpLag だ。CatchUp を呼ぶまで、読みモデルは書き込みに遅れ、その遅れの件数を Lag が返す。テストで、注文して支払った直後、CatchUp する前は読みモデルがまだ空で、CatchUp すると追いつくことを固定した。実際のシステムでは、読み取り側は書き込み側のイベントを非同期に購読するので、この遅れはミリ秒単位でも必ず存在する。「注文した直後に一覧を見たら、まだ載っていない」ということが起こりうる。これが結果整合(eventual consistency)だ。強い整合性(書いた瞬間に読める)を諦める代わりに、読みと書きを独立にスケールでき、読みモデルを自由に増やせる。この取引をどう受け入れるかが CQRS を使う判断になる。銀行の残高のように即時一貫性が要る箇所には向かず、注文履歴やダッシュボードのように少しの遅れが許される箇所に向く。

動かす

下のデモは、注文・支払い・取消のコマンドを送りながら、書き込み側のイベントログと、そこから導かれる 3 つの読みモデルを見る。CatchUp する前は読みモデルが書き込みに遅れ(結果整合)、CatchUp で追いつく様子も確かめられる。

デモCQRS(注文)読みが 1 遅れ
書き込み側(イベントログ)
v1placed o1
v2paid o1
v3placed o2未処理
読み取り側(3 つの読みモデル)
① 注文状況
② 顧客ごとの購入額
alice:1000
③ 総売上
1000

書き込み側はコマンドを検証してイベントを流す。読み取り側は同じイベント列を畳んで、注文状況・顧客ごとの 購入額・総売上という 3 つの読みモデルを同時に作る。1 つのイベントから用途別のビューが何個でも組める。 CatchUp する前は読みモデルが書き込みに遅れる(未処理・結果整合)。CatchUp で追いつく。読みと書きを 分けることで、それぞれを独立に最適化できる代わりに、この遅れを受け入れる。

設計の観点

  • 分離の判断は整合性の要求で: 即時一貫性が要る箇所(残高・在庫の確定)には向かない。少しの遅れが許される読み(一覧・集計・ダッシュボード)に向く
  • 読みモデルは使い捨て: イベントから作り直せるので、スキーマ変更は読みモデルを捨てて再構築すればいい。書き込み側は影響を受けない
  • 書きと読みを別々にスケール: 読みが重いなら読みモデルだけ増やす。書きが重いなら書き込み側だけ強化する。負荷特性に応じて独立に
  • イベントソーシングとは別概念: CQRS は読み書きの分離、イベントソーシングはイベントを真実とする保存。相性はよいが、どちらか一方だけでも使える
  • 複雑さの代償: 分離・非同期・結果整合はシステムを複雑にする。単純な CRUD で足りるなら、CQRS は過剰。適用範囲を絞る

対照と実例

単一モデル(CRUD)CQRS
読み書きモデル共通分離
読みビュー1 つ(書きモデル由来)用途別に複数
整合性強い(即時)結果整合(遅れる)
スケール一体読み書き独立
複雑さ低い高い

裏どり:

  • Martin Fowler: CQRS: 概念の定番の解説。適用すべき箇所とすべきでない箇所の整理
  • Greg Young の CQRS 講演: イベントソーシングと組み合わせた CQRS の原典的な議論
  • 読み書き分離レプリカ: DB のリードレプリカも、粗い意味では読み書き分離。結果整合の遅れは同型の問題
  • Elasticsearch を読みモデルに: 書き込みは RDB、検索用の読みモデルは Elasticsearch、という CQRS の実運用パターン

簡略化したこと

  • 同期の CatchUp: 実物は読み側が非同期にイベントを購読する。ここは明示的に呼んで遅れを見せる
  • 単一プロセス: 書きと読みが同じメモリ。実物は別サービス・別 DB(書きは正規化、読みは非正規化)
  • 単純なドメイン: 注文の 3 状態のみ。実物の不変条件はもっと複雑
  • 並行制御なし: 同時コマンドの競合や順序保証は扱わない

参考資料