mybook

イベントソーシング — 状態ではなく事実を記録する

「なぜ注文が500円引きになったか分からない」

本番稼働2週間後。カオリに深刻な問い合わせが来た。

「注文#1234の合計金額がおかしい。内訳を調べたいが、履歴が残っていない」

調べると、注文作成後に管理者が金額を修正していた。しかしデータベースには最終的な状態しか残っていない。「誰が、いつ、なぜ修正したか」の証跡が全くない。

通常のデータベースは現在の状態しか保存しない。UPDATE orders SET total_price = 4500 WHERE id = 1234 を実行すると、元の値は消える。

INFO

イベントソーシングは、システムの状態変化を「イベント(事実)」として永続化します。現在の状態は、全てのイベントを再生(replay)することで導出します。銀行の通帳や会計の仕訳帳と同じ発想です。

イベントソーシングの基本概念

Loading diagram...

銀行口座で例えると:

  • 通常のDB: 残高 = 50,000円 (現在の状態のみ)
  • イベントソーシング: 入金10,000→出金5,000→入金30,000→出金…(全履歴)

イベントの設計

# app/events/base_event.rb
class BaseEvent
  attr_reader :event_id, :occurred_at, :aggregate_id, :aggregate_type, :version
  
  def initialize(aggregate_id:, version:, **data)
    @event_id = SecureRandom.uuid
    @aggregate_id = aggregate_id
    @aggregate_type = self.class.aggregate_type
    @version = version
    @occurred_at = Time.current
    @data = data
  end
  
  def to_h
    {
      event_id: event_id,
      event_type: self.class.event_type,
      aggregate_id: aggregate_id,
      aggregate_type: aggregate_type,
      version: version,
      occurred_at: occurred_at,
      data: @data
    }
  end
end
 
# app/events/order_events.rb
class OrderPlaced < BaseEvent
  def self.event_type = 'order.placed'
  def self.aggregate_type = 'Order'
  
  attr_reader :user_id, :product_id, :quantity, :unit_price
  
  def initialize(aggregate_id:, version:, user_id:, product_id:, quantity:, unit_price:)
    super(aggregate_id: aggregate_id, version: version,
          user_id: user_id, product_id: product_id,
          quantity: quantity, unit_price: unit_price)
    @user_id = user_id
    @product_id = product_id
    @quantity = quantity
    @unit_price = unit_price
  end
end
 
class OrderConfirmed < BaseEvent
  def self.event_type = 'order.confirmed'
  def self.aggregate_type = 'Order'
  
  attr_reader :confirmed_by_user_id
  
  def initialize(aggregate_id:, version:, confirmed_by_user_id:)
    super(aggregate_id: aggregate_id, version: version,
          confirmed_by_user_id: confirmed_by_user_id)
    @confirmed_by_user_id = confirmed_by_user_id
  end
end
 
class OrderPriceAdjusted < BaseEvent
  def self.event_type = 'order.price_adjusted'
  def self.aggregate_type = 'Order'
  
  attr_reader :original_price, :adjusted_price, :reason, :adjusted_by_user_id
  
  def initialize(aggregate_id:, version:, original_price:, adjusted_price:, reason:, adjusted_by_user_id:)
    super(aggregate_id: aggregate_id, version: version,
          original_price: original_price, adjusted_price: adjusted_price,
          reason: reason, adjusted_by_user_id: adjusted_by_user_id)
    @original_price = original_price
    @adjusted_price = adjusted_price
    @reason = reason
    @adjusted_by_user_id = adjusted_by_user_id
  end
end

イベントストア

# app/event_store/event_store.rb
class EventStore
  # イベントの追加(追記のみ、更新・削除なし)
  def append(stream_name, events, expected_version: nil)
    events = Array(events)
    
    ActiveRecord::Base.transaction do
      if expected_version
        current_version = EventRecord.where(stream_name: stream_name).maximum(:version) || -1
        raise OptimisticLockError if current_version != expected_version
      end
      
      events.each do |event|
        EventRecord.create!(
          stream_name: stream_name,
          event_id: event.event_id,
          event_type: event.class.event_type,
          aggregate_id: event.aggregate_id,
          aggregate_type: event.aggregate_type,
          version: event.version,
          data: event.to_h[:data].to_json,
          occurred_at: event.occurred_at
        )
      end
      
      # イベントの発行(サブスクライバーへの通知)
      events.each { |e| publish_to_subscribers(e) }
    end
  end
  
  # ストリームからイベントを読み取る
  def read_stream(stream_name, from_version: 0)
    EventRecord
      .where(stream_name: stream_name, version: from_version..)
      .order(:version)
      .map { |r| deserialize(r) }
  end
  
  # 任意の時点までのイベントを読み取る(タイムトラベル!)
  def read_stream_until(stream_name, timestamp)
    EventRecord
      .where(stream_name: stream_name)
      .where('occurred_at <= ?', timestamp)
      .order(:version)
      .map { |r| deserialize(r) }
  end
  
  private
  
  def deserialize(record)
    event_class = EVENT_TYPES.fetch(record.event_type)
    data = JSON.parse(record.data, symbolize_names: true)
    event_class.new(
      aggregate_id: record.aggregate_id,
      version: record.version,
      **data
    )
  end
  
  EVENT_TYPES = {
    'order.placed' => OrderPlaced,
    'order.confirmed' => OrderConfirmed,
    'order.price_adjusted' => OrderPriceAdjusted
  }.freeze
end

イベントストアのスキーマ

# db/migrate/create_event_records.rb
class CreateEventRecords < ActiveRecord::Migration[7.1]
  def change
    create_table :event_records, id: :uuid do |t|
      t.string :stream_name, null: false
      t.uuid :event_id, null: false
      t.string :event_type, null: false
      t.uuid :aggregate_id, null: false
      t.string :aggregate_type, null: false
      t.integer :version, null: false
      t.jsonb :data, null: false, default: {}
      t.datetime :occurred_at, null: false
      
      t.timestamps
    end
    
    add_index :event_records, :event_id, unique: true
    add_index :event_records, [:stream_name, :version], unique: true
    add_index :event_records, [:aggregate_type, :aggregate_id]
    add_index :event_records, :occurred_at
    add_index :event_records, :event_type
  end
end

集約(Aggregate)の実装

# app/aggregates/order_aggregate.rb
class OrderAggregate
  attr_reader :id, :user_id, :product_id, :quantity, :unit_price, :status, :version
  
  def initialize
    @version = -1
    @uncommitted_events = []
  end
  
  # クラスメソッド:イベントから集約を再構築する
  def self.from_events(events)
    aggregate = new
    events.each { |event| aggregate.apply(event, record: false) }
    aggregate
  end
  
  # コマンド: 注文を配置する
  def place(user_id:, product_id:, quantity:, unit_price:)
    raise "既に配置済み" if @id
    
    event = OrderPlaced.new(
      aggregate_id: SecureRandom.uuid,
      version: @version + 1,
      user_id: user_id,
      product_id: product_id,
      quantity: quantity,
      unit_price: unit_price
    )
    
    apply(event)
  end
  
  # コマンド: 注文を確定する
  def confirm(confirmed_by:)
    raise "注文は pending 状態でなければなりません" unless @status == 'pending'
    
    event = OrderConfirmed.new(
      aggregate_id: @id,
      version: @version + 1,
      confirmed_by_user_id: confirmed_by
    )
    
    apply(event)
  end
  
  # コマンド: 価格を調整する(監査ログが重要)
  def adjust_price(new_price:, reason:, adjusted_by:)
    event = OrderPriceAdjusted.new(
      aggregate_id: @id,
      version: @version + 1,
      original_price: @unit_price * @quantity,
      adjusted_price: new_price,
      reason: reason,
      adjusted_by_user_id: adjusted_by
    )
    
    apply(event)
  end
  
  def uncommitted_events
    @uncommitted_events.dup
  end
  
  def mark_changes_as_committed
    @uncommitted_events.clear
  end
  
  private
  
  def apply(event, record: true)
    # イベントを状態に反映する
    case event
    when OrderPlaced
      @id = event.aggregate_id
      @user_id = event.user_id
      @product_id = event.product_id
      @quantity = event.quantity
      @unit_price = event.unit_price
      @status = 'pending'
    when OrderConfirmed
      @status = 'confirmed'
    when OrderPriceAdjusted
      @adjusted_price = event.adjusted_price
    end
    
    @version = event.version
    @uncommitted_events << event if record
  end
end

リポジトリ(イベントソーシング版)

# app/repositories/order_event_repository.rb
class OrderEventRepository
  def initialize(event_store)
    @event_store = event_store
  end
  
  def save(aggregate)
    stream_name = stream_for(aggregate.id)
    events = aggregate.uncommitted_events
    
    @event_store.append(
      stream_name,
      events,
      expected_version: aggregate.version - events.size
    )
    
    aggregate.mark_changes_as_committed
    aggregate
  end
  
  def find(aggregate_id)
    stream_name = "Order-#{aggregate_id}"
    events = @event_store.read_stream(stream_name)
    
    raise ActiveRecord::RecordNotFound if events.empty?
    
    OrderAggregate.from_events(events)
  end
  
  # タイムトラベル機能!
  def find_at(aggregate_id, timestamp)
    stream_name = "Order-#{aggregate_id}"
    events = @event_store.read_stream_until(stream_name, timestamp)
    
    raise ActiveRecord::RecordNotFound if events.empty?
    
    OrderAggregate.from_events(events)
  end
  
  private
  
  def stream_for(id) = "Order-#{id}"
end

タイムトラベル機能の使い方

# 「昨日の17:00時点での注文状態」を確認できる
repository = OrderEventRepository.new(event_store)
order_yesterday = repository.find_at('order-1234', 1.day.ago)
 
puts order_yesterday.status      # => 'pending'
puts order_yesterday.unit_price  # => 2000
 
# 現在の状態
order_now = repository.find('order-1234')
puts order_now.status            # => 'confirmed'
puts order_now.adjusted_price    # => 1500

「誰が、いつ、なぜ500円引きにしたか」が全部分かる。

AWSでのイベントソーシング構成

Loading diagram...
# EventBridge + Lambda でプロジェクターを構築
event_bridge_rule:
  EventPattern:
    source: ["myapp.orders"]
    detail-type: 
      - "order.placed"
      - "order.confirmed"
      - "order.price_adjusted"
 
lambda_projector:
  handler: projector.handler
  events:
    - eventBridge:
        eventBus: myapp-events
        pattern:
          source:
            - myapp.orders
  environment:
    READ_DB_URL: ${ssm:/myapp/read_db_url}

スナップショット

イベントが蓄積すると、リプレイに時間がかかる。スナップショットで対策する。

# app/aggregates/snapshotting_order_repository.rb
class SnapshottingOrderRepository
  SNAPSHOT_INTERVAL = 50  # 50イベントごとにスナップショット
  
  def find(aggregate_id)
    snapshot = SnapshotRecord.find_by(aggregate_id: aggregate_id)
    
    if snapshot
      # スナップショット以降のイベントのみリプレイ
      aggregate = restore_from_snapshot(snapshot)
      newer_events = @event_store.read_stream(
        "Order-#{aggregate_id}",
        from_version: snapshot.version + 1
      )
      newer_events.each { |e| aggregate.apply_event(e) }
      aggregate
    else
      # 全イベントをリプレイ
      events = @event_store.read_stream("Order-#{aggregate_id}")
      OrderAggregate.from_events(events)
    end
  end
  
  def save(aggregate)
    super
    
    # スナップショットを作成するタイミング
    if aggregate.version % SNAPSHOT_INTERVAL == 0
      SnapshotRecord.create_or_update!(
        aggregate_id: aggregate.id,
        aggregate_type: 'Order',
        version: aggregate.version,
        state: aggregate.to_snapshot.to_json
      )
    end
  end
end

イベントソーシングのトレードオフ

メリット:
✅ 完全な監査ログ(誰が、いつ、何をしたか)
✅ タイムトラベル(任意の時点の状態を再現)
✅ バグの原因追跡が容易
✅ 読み取りモデルを後から追加できる(過去から再構築)
✅ CQRSとの相性が抜群

デメリット:
❌ 実装の複雑さが増す
❌ パフォーマンス(リプレイのコスト)
❌ スキーマ変更(イベントの後方互換性)
❌ デバッグが難しい
❌ 学習コストが高い

WARNING

イベントソーシングは、監査ログが法的・業務上の要件として必要、または状態の変化の追跡が本質的に重要なドメイン(金融、医療、法律)に最適です。全てのシステムに適用するのは過剰設計です。

カオリは決断した。「注文の金額調整と、在庫の変動だけイベントソーシングを適用する。残りは通常のActiveRecordで十分だ」

選択的な適用。それがカオリの現実的な答えだった。


次章ではマクロな視点に戻る。「サービス指向アーキテクチャ(SOA)」の歴史と教訓、そしてなぜマイクロサービスが生まれたのかを学ぶ。