CQRS
実装:
messaging/cqrs// 実行:go test ./messaging/cqrs/
普通のアプリは1つのモデルで読み書きの両方をこなす。だが読みと書きは要求が違う。書きは正しさが命で、読みは速さと使いやすい形が命だ。両立させると両方が中途半端になる。CQRSはこれを分ける。書き込み側は状態を検証してイベントを流し、読み取り側はそのイベントを畳んで用途別のビューを作る。1つのイベント列から注文状況・顧客ごとの購入額・総売上を独立に組める。分離を実装し、代償の結果整合まで確かめる。
この章で作るもの
イベントソーシングで、状態を保存せずイベントから導く仕組みを作った。CQRS(Command Query Responsibility Segregation)は、その導き方を「書き込み」と「読み取り」で分ける考え方だ。多くのアプリは、1 つのデータモデルで読み書きの両方をこなす。同じテーブル、同じオブジェクトに書き、そこから読む。だが読みと書きは、求めるものがまるで違う。
書き込みが大事にするのは正しさだ。「支払済みの注文は取り消せない」「在庫を超える注文は受けない」といった不変条件を守り、状態を矛盾なく変える。一方、読み取りが大事にするのは速さと形だ。「この顧客の購入履歴を一覧で」「今日の売上合計を」といった問いに、素早く、都合のよい形で答えたい。1 つのモデルで両方を最適化しようとすると、書きのための正規化された構造が読みには使いにくく、読みのための非正規化が書きの整合性を危うくする。どちらも中途半端になる。CQRS は、書き込みモデルと読み取りモデルを分けてこれを解く。書き込み側はイベントを流し、読み取り側はそのイベントから、用途ごとに最適な読みモデル(射影)を組む。
コマンド ──▶ 書き込み側 ──イベント──▶ 読み取り側 ┬─▶ 注文状況の一覧
(検証・整合性) (追記ログ) ├─▶ 顧客ごとの購入額
└─▶ 総売上
1 つのイベント列から、用途別のビューを何個でも順に見ていく。
- 書きと読みを分ける: 書き込み側は正しさ、読み取り側は速さと形。それぞれを独立に最適化できる
- 1 イベント列から複数ビュー: 同じイベントを別々に畳んで、用途ごとの読みモデルをいくつも作る
- 代償は結果整合: 読みモデルは書きに少し遅れて追いつく。その間は古い値を返す
① 書き込み側: コマンドを検証してイベントを流す
まず書き込み側だ。コマンド(注文する、支払う、取り消す)を受け、状態を検証してからイベントを流す。検証のための状態は、イベントを畳んで導く(イベントソーシングと同じ):
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 つのイベント列から、何種類でも:
// 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 つの読みモデルが正しく構築されることを固定した。
③ 結果整合: 読みは書きに遅れる
分離には代償がある。読みモデルは、書き込みと同じ瞬間には更新されない。書き込み側がイベントを流し、読み取り側がそれを処理して初めて、読みモデルに反映される。この処理が追いつくまで、読みモデルは少し古い値を返す。それを担うのが、上のコードにある CatchUp と Lag だ。CatchUp を呼ぶまで、読みモデルは書き込みに遅れ、その遅れの件数を Lag が返す。テストで、注文して支払った直後、CatchUp する前は読みモデルがまだ空で、CatchUp すると追いつくことを固定した。実際のシステムでは、読み取り側は書き込み側のイベントを非同期に購読するので、この遅れはミリ秒単位でも必ず存在する。「注文した直後に一覧を見たら、まだ載っていない」ということが起こりうる。これが結果整合(eventual consistency)だ。強い整合性(書いた瞬間に読める)を諦める代わりに、読みと書きを独立にスケールでき、読みモデルを自由に増やせる。この取引をどう受け入れるかが CQRS を使う判断になる。銀行の残高のように即時一貫性が要る箇所には向かず、注文履歴やダッシュボードのように少しの遅れが許される箇所に向く。
動かす
下のデモは、注文・支払い・取消のコマンドを送りながら、書き込み側のイベントログと、そこから導かれる 3 つの読みモデルを見る。CatchUp する前は読みモデルが書き込みに遅れ(結果整合)、CatchUp で追いつく様子も確かめられる。
書き込み側はコマンドを検証してイベントを流す。読み取り側は同じイベント列を畳んで、注文状況・顧客ごとの 購入額・総売上という 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 状態のみ。実物の不変条件はもっと複雑
- 並行制御なし: 同時コマンドの競合や順序保証は扱わない
参考資料
- Martin Fowler: CQRS — 概念と適用判断
- Microsoft: CQRS pattern — 実装パターンと利点・欠点
- Greg Young: CQRS Documents — 原典的な詳説
- 実装: messaging/cqrs