mybook

ドメインイベント — 変化を通知する

「注文確定したら、在庫を引き当てて、メールを送って、ポイントを付与して、分析データを送って、Slackに通知して……」

リナは要件一覧を見ながらため息をついた。注文確定という1つのアクションに対して、後続処理が6種類もある。しかも要件は毎月増える。「プッシュ通知も追加してほしい」「注文確定で返金ポイントも計算してほしい」...

これを全部 Order#confirm! の中に書いたら、また Fat Model の悪夢だ。かといって、アプリケーションサービスに書いても、サービスが肥大化する。

# アンチパターン: confirm! に後続処理を全部書く
class Order
  def confirm!
    update!(status: 'confirmed')
 
    # 在庫引き当て
    order_items.each { |item| Stock.reserve!(item) }
 
    # メール送信
    OrderConfirmationMailer.send_confirmation(self).deliver_later
 
    # ポイント付与
    LoyaltyService.add_points(customer, total_amount)
 
    # 分析データ
    Analytics.track('order_confirmed', order_id: id, amount: total_amount)
 
    # Slack通知
    SlackNotifier.notify_ops("新規注文: #{id}")
 
    # プッシュ通知
    PushNotificationService.send(customer, 'ご注文ありがとうございます')
 
    # 在庫アラート確認
    InventoryAlertService.check_low_stock(order_items.map(&:product))
  end
end

Order が注文管理以外の責務を7つも持っている。新しい後続処理が追加されるたびに Order クラスを変更しなければならない。これはOpen-Closed原則(拡張に開いている、修正に閉じている)の違反だ。

ドメインイベントとは

ドメインイベント(Domain Event)は、ドメイン内で起きた「何らかの出来事」を表すオブジェクトだ。

過去形で命名する——「何かが起きた」という事実を表すからだ。

  • OrderConfirmed(注文が確定した)
  • StockReserved(在庫が引き当てられた)
  • PaymentProcessed(決済が処理された)
  • OrderCancelled(注文がキャンセルされた)
Loading diagram...

INFO

ドメインイベントを使うと、Order#confirm! は「注文を確定した」という事実を記録するだけでよくなる。後続処理は疎結合なハンドラーが担い、Orderの責務が純粋になる。新しい後続処理は Order を変えずにハンドラーを追加するだけで実装できる。

イベントオブジェクトの実装

# ドメインイベントの基底クラス
module DomainEvents
  class DomainEvent
    attr_reader :event_id, :occurred_at, :metadata
 
    def initialize(metadata: {})
      @event_id = SecureRandom.uuid
      @occurred_at = Time.current.utc  # UTCで統一
      @metadata = metadata.freeze
      freeze
    end
 
    def event_type
      # 例: "order_context.order_confirmed"
      self.class.name
          .underscore
          .gsub('/', '.')
          .delete_prefix('domain_events.')
    end
 
    def to_h
      {
        event_id: event_id,
        event_type: event_type,
        occurred_at: occurred_at.iso8601,
        payload: payload,
        metadata: metadata
      }
    end
 
    def to_json(*args)
      to_h.to_json(*args)
    end
 
    private
 
    def payload
      {}  # サブクラスでオーバーライドする
    end
  end
end
 
# 注文確定イベント
module OrderContext
  class OrderConfirmed < DomainEvents::DomainEvent
    attr_reader :order_id, :customer_id, :total_amount, :order_items, :confirmed_at
 
    def initialize(order_id:, customer_id:, total_amount:, order_items:, confirmed_at:, **options)
      super(**options)
      @order_id = order_id
      @customer_id = customer_id
      @total_amount = total_amount   # SharedKernel::Money
      @order_items = order_items.freeze
      @confirmed_at = confirmed_at
    end
 
    private
 
    def payload
      {
        order_id: order_id,
        customer_id: customer_id,
        total_amount: total_amount.to_h,
        order_items: order_items,
        confirmed_at: confirmed_at.iso8601
      }
    end
  end
 
  class OrderCancelled < DomainEvents::DomainEvent
    attr_reader :order_id, :customer_id, :reason, :cancelled_at
 
    def initialize(order_id:, customer_id:, reason:, cancelled_at:, **options)
      super(**options)
      @order_id = order_id
      @customer_id = customer_id
      @reason = reason
      @cancelled_at = cancelled_at
    end
 
    private
 
    def payload
      {
        order_id: order_id,
        customer_id: customer_id,
        reason: reason,
        cancelled_at: cancelled_at.iso8601
      }
    end
  end
end

ドメインモデルへの組み込み

集約がイベントを「積む」。保存後に発行するためにキューを使う。

module OrderContext
  class Order
    def confirm!
      raise InvalidStateTransition unless can_confirm?
      raise EmptyOrder if order_items.empty?
 
      @status = OrderStatus::CONFIRMED
      @confirmed_at = Time.current
 
      # イベントをキューに積む(この時点ではまだ発行しない)
      # 保存成功を確認してから発行するため
      record_event(OrderConfirmed.new(
        order_id: id,
        customer_id: customer_id,
        total_amount: total_amount,
        order_items: order_items.map(&:to_event_payload),
        confirmed_at: @confirmed_at
      ))
 
      self
    end
 
    def cancel!(reason:)
      raise InvalidStateTransition unless can_cancel?
 
      @status = OrderStatus::CANCELLED
      @cancellation_reason = reason
      @cancelled_at = Time.current
 
      record_event(OrderCancelled.new(
        order_id: id,
        customer_id: customer_id,
        reason: reason,
        cancelled_at: @cancelled_at
      ))
 
      self
    end
 
    def domain_events
      @domain_events.dup.freeze
    end
 
    def clear_events
      @domain_events = []
      self
    end
 
    private
 
    def record_event(event)
      @domain_events ||= []
      @domain_events << event
    end
  end
end

イベントバスの実装

Rails内での同期処理(ActiveSupport::Notifications)

モノリスの中での疎結合なイベント処理。

# config/initializers/event_bus.rb
module EventBus
  module_function
 
  def publish(event)
    Rails.logger.tagged('EventBus') do
      Rails.logger.info("Publish: #{event.event_type} [#{event.event_id}]")
    end
 
    ActiveSupport::Notifications.instrument(
      "domain_event.#{event.event_type}",
      event: event
    )
  rescue => e
    # イベント発行エラーはログに残すが、呼び出し元の処理は止めない
    Rails.logger.error("EventBus publish failed: #{e.message}")
    ErrorTracker.capture(e, event_type: event.event_type)
  end
 
  def subscribe(event_type_pattern, &handler)
    ActiveSupport::Notifications.subscribe(
      "domain_event.#{event_type_pattern}"
    ) do |name, start, finish, id, payload|
      event = payload[:event]
      begin
        handler.call(event)
      rescue => e
        Rails.logger.error(
          "EventHandler failed for #{event.event_type}: #{e.message}"
        )
        ErrorTracker.capture(e, event_id: event.event_id)
        raise  # エラーを再発生させてハンドラーの失敗を明示
      end
    end
  end
end
 
# config/initializers/event_handlers.rb
# 各ハンドラーをイベントタイプに紐付ける
 
# 注文確定時の後続処理
EventBus.subscribe('order_context.order_confirmed') do |event|
  StockReservationHandler.new.call(event)
end
 
EventBus.subscribe('order_context.order_confirmed') do |event|
  OrderConfirmationEmailHandler.new.call(event)
end
 
EventBus.subscribe('order_context.order_confirmed') do |event|
  LoyaltyPointsHandler.new.call(event)
end
 
EventBus.subscribe('order_context.order_confirmed') do |event|
  OrderAnalyticsHandler.new.call(event)
end
 
# 注文キャンセル時の後続処理
EventBus.subscribe('order_context.order_cancelled') do |event|
  StockReleaseHandler.new.call(event)
end
 
EventBus.subscribe('order_context.order_cancelled') do |event|
  RefundProcessingHandler.new.call(event)
end

ハンドラーの実装

各ハンドラーは単一の責務を持つ小さなクラス。

# 在庫引き当てハンドラー
class StockReservationHandler
  def initialize(
    stock_repository: InventoryContext::ActiveRecordStockRepository.new,
    order_repository: OrderContext::ActiveRecordOrderRepository.new
  )
    @stock_repository = stock_repository
    @order_repository = order_repository
    @reservation_service = InventoryContext::StockReservationService.new(
      stock_repository: @stock_repository
    )
  end
 
  def call(event)
    # イベントのデータで在庫を引き当てる
    event.order_items.each do |item|
      @reservation_service.reserve(
        product_id: item[:product_id],
        quantity: item[:quantity],
        order_id: event.order_id
      )
    end
  rescue InventoryContext::InsufficientStock => e
    # 在庫不足 → 注文をキャンセルする(補償トランザクション)
    Rails.logger.warn("在庫不足でキャンセル: order=#{event.order_id}, #{e.message}")
    cancel_order_due_to_stock_shortage(event.order_id, e.message)
  end
 
  private
 
  def cancel_order_due_to_stock_shortage(order_id, reason)
    order = @order_repository.find(order_id)
    order.cancel!(reason: "在庫不足のため自動キャンセル: #{reason}")
    @order_repository.save(order)
    EventBus.publish(order.domain_events.last)
    order.clear_events
  end
end
 
# メール送信ハンドラー
class OrderConfirmationEmailHandler
  def call(event)
    OrderConfirmationMailer
      .confirmation(order_id: event.order_id, customer_id: event.customer_id)
      .deliver_later(queue: :transactional_email)
  rescue => e
    # メール送信失敗は注文自体には影響させない
    Rails.logger.error("確認メール送信失敗: order=#{event.order_id}, #{e.message}")
    ErrorTracker.capture(e, order_id: event.order_id)
  end
end
 
# ポイント付与ハンドラー
class LoyaltyPointsHandler
  POINTS_PER_YEN = Rational(1, 100)  # 100円 = 1ポイント
 
  def initialize(loyalty_service: LoyaltyService.new)
    @loyalty_service = loyalty_service
  end
 
  def call(event)
    points_to_add = (event.total_amount[:amount] * POINTS_PER_YEN).to_i
    return if points_to_add <= 0
 
    @loyalty_service.add_points(
      customer_id: event.customer_id,
      points: points_to_add,
      reason: "注文確定: #{event.order_id}",
      reference_id: event.order_id
    )
  end
end
 
# 分析データハンドラー
class OrderAnalyticsHandler
  def call(event)
    AnalyticsService.track(
      event_name: 'order_confirmed',
      user_id: event.customer_id,
      properties: {
        order_id: event.order_id,
        total_amount: event.total_amount[:amount],
        item_count: event.order_items.sum { |i| i[:quantity] },
        confirmed_at: event.confirmed_at
      }
    )
  rescue => e
    # 分析失敗はサービスに影響させない
    Rails.logger.error("Analytics失敗: #{e.message}")
  end
end

アプリケーションサービスでのイベント発行

イベント発行のタイミングが重要だ。

class ConfirmOrderUseCase
  def call(order_id:, customer_id:)
    order = @order_repository.find(order_id)
    raise Unauthorized unless order.customer_id == customer_id
 
    # 1. ドメインオブジェクトへの操作(イベントをキューに積む)
    order.confirm!
 
    # 2. トランザクション内で: Order保存 + Outboxにイベント記録
    ApplicationRecord.transaction do
      @order_repository.save(order)
      record_outbox_events(order.domain_events)  # 同じトランザクションで
    end
 
    # 3. トランザクション成功後: イベントを発行
    # (失敗した場合、Outboxからリトライできる)
    publish_events(order.domain_events)
    order.clear_events
 
    ConfirmOrderResult.success(order: order)
  end
 
  private
 
  def record_outbox_events(events)
    events.each do |event|
      EventOutbox.create!(
        event_id: event.event_id,
        event_type: event.event_type,
        payload: event.to_json,
        status: 'pending',
        created_at: Time.current
      )
    end
  end
 
  def publish_events(events)
    events.each { |event| EventBus.publish(event) }
  end
end

WARNING

イベント発行のタイミングは重要だ。データベース保存の後にイベントを発行しないと、「イベントは発行されたが保存が失敗した」という矛盾状態が発生する。Outboxパターンで「保存とイベント記録を同一トランザクションで行う」ことで原子性を保つ。

Outboxパターン — 信頼性の高いイベント発行

# db/migrate/20240115_create_event_outboxes.rb
class CreateEventOutboxes < ActiveRecord::Migration[7.1]
  def change
    create_table :event_outboxes do |t|
      t.string :event_id, null: false
      t.string :event_type, null: false
      t.text :payload, null: false
      t.string :status, default: 'pending', null: false  # pending / published / failed
      t.integer :retry_count, default: 0
      t.datetime :published_at
      t.datetime :failed_at
      t.text :error_message
 
      t.timestamps
    end
 
    add_index :event_outboxes, :status
    add_index :event_outboxes, :event_id, unique: true
    add_index :event_outboxes, [:status, :created_at]
  end
end
 
# バックグラウンドジョブでOutboxからイベントを発行
class OutboxEventPublisherJob < ApplicationJob
  queue_as :default
  sidekiq_options retry: 3
 
  def perform
    EventOutbox.where(status: 'pending')
               .where('retry_count < 5')
               .order(:created_at)
               .limit(100)
               .each do |outbox|
      begin
        payload = JSON.parse(outbox.payload, symbolize_names: true)
        EventBus.publish_raw(
          event_type: outbox.event_type,
          payload: payload
        )
        outbox.update!(status: 'published', published_at: Time.current)
      rescue => e
        outbox.update!(
          retry_count: outbox.retry_count + 1,
          error_message: e.message,
          status: outbox.retry_count >= 4 ? 'failed' : 'pending'
        )
      end
    end
  end
end

AWS EventBridgeによる非同期処理

FreshCartが成長してマイクロサービス化する場合。

Loading diagram...
# AWS EventBridgeへのイベント送信
class EventBridgeAdapter
  def initialize(
    client: Aws::EventBridge::Client.new(region: ENV['AWS_REGION']),
    event_bus_name: ENV['EVENT_BUS_NAME']
  )
    @client = client
    @event_bus_name = event_bus_name
  end
 
  def publish(event)
    response = @client.put_events(
      entries: [{
        source: "freshcart.#{event.event_type.split('.').first}",
        detail_type: event.event_type,
        detail: event.to_json,
        event_bus_name: @event_bus_name,
        time: event.occurred_at
      }]
    )
 
    failed = response.entries.select { |e| e.error_code }
    if failed.any?
      raise "EventBridge発行失敗: #{failed.map(&:error_message).join(', ')}"
    end
 
    Rails.logger.info("EventBridge published: #{event.event_type}")
  rescue Aws::EventBridge::Errors::ServiceError => e
    Rails.logger.error("EventBridge error: #{e.message}")
    raise
  end
end
 
# Lambda関数(在庫サービス側)でのイベント受信
# handler.rb (Golang風にRubyで記述)
module InventoryServiceHandler
  def self.handle_order_confirmed(event_detail)
    order_id = event_detail['order_id']
    order_items = event_detail['order_items']
 
    stock_service = InventoryContext::StockReservationService.new(
      stock_repository: InventoryContext::DynamoDBStockRepository.new
    )
 
    order_items.each do |item|
      stock_service.reserve(
        product_id: item['product_id'],
        quantity: item['quantity'],
        order_id: order_id
      )
    end
 
    { statusCode: 200, body: 'OK' }
  rescue InventoryContext::InsufficientStock => e
    # 在庫不足: 注文サービスに補償イベントを送信
    EventBridgeAdapter.new.publish(
      StockReservationFailed.new(
        order_id: order_id,
        reason: e.message
      )
    )
    { statusCode: 200, body: 'Handled' }
  end
end

イベントソーシング(発展的な概念)

イベントの履歴から状態を再構築するアプローチ。

# すべての状態変化をイベントとして永続化
class EventSourcedOrder
  attr_reader :id, :status, :order_items, :total_amount
 
  def self.create_from_events(events)
    order = allocate
    order.instance_variable_set(:@domain_events_applied, [])
    events.each { |event| order.apply(event) }
    order
  end
 
  def apply(event)
    case event
    when OrderCreated
      @id = event.order_id
      @customer_id = event.customer_id
      @status = :pending
      @order_items = []
    when OrderItemAdded
      @order_items << { product_id: event.product_id, quantity: event.quantity }
    when OrderConfirmed
      @status = :confirmed
      @confirmed_at = event.confirmed_at
    when OrderCancelled
      @status = :cancelled
      @cancellation_reason = event.reason
    end
    @domain_events_applied << event
    self
  end
 
  # 任意の時点の状態を再現できる
  def self.state_at(order_id:, point_in_time:)
    events = EventStore.find_by_aggregate_id(order_id)
                       .where('occurred_at <= ?', point_in_time)
    create_from_events(events)
  end
end

INFO

イベントソーシングは強力だが複雑さも高い。「なぜこの状態になったか」の完全な履歴を持てるため、監査要件の厳しい金融系に特に有効だ。FreshCartのような一般的なECサイトでは、通常のCRUD + ドメインイベントの組み合わせが現実的だ。

テスト戦略

RSpec.describe OrderContext::Order, 'ドメインイベント' do
  let(:order) do
    order = OrderContext::Order.new(
      id: 'order-1',
      customer_id: 'customer-1',
      delivery_address: build_address
    )
    order.add_item(
      product_id: 'product-1',
      product_name: 'りんご',
      unit_price: SharedKernel::Money.new(amount: 500, currency: :jpy),
      quantity: 2
    )
    order
  end
 
  describe '注文確定後のドメインイベント' do
    before { order.confirm! }
 
    it 'OrderConfirmedイベントが1つ発行される' do
      expect(order.domain_events.size).to eq(1)
      expect(order.domain_events.first).to be_a(OrderContext::OrderConfirmed)
    end
 
    it 'OrderConfirmedイベントに正しいデータが含まれる' do
      event = order.domain_events.first
 
      expect(event.order_id).to eq('order-1')
      expect(event.customer_id).to eq('customer-1')
      expect(event.total_amount).to eq(SharedKernel::Money.new(amount: 1000, currency: :jpy))
    end
 
    it 'clear_eventsでイベントがクリアされる' do
      order.clear_events
      expect(order.domain_events).to be_empty
    end
  end
end
 
RSpec.describe StockReservationHandler do
  let(:stock_repository) { InventoryContext::InMemoryStockRepository.new }
  let(:handler) { described_class.new(stock_repository: stock_repository) }
  let(:event) do
    OrderContext::OrderConfirmed.new(
      order_id: 'order-1',
      customer_id: 'customer-1',
      total_amount: SharedKernel::Money.new(amount: 1000, currency: :jpy),
      order_items: [
        { product_id: 'product-1', quantity: 2, product_name: 'りんご', unit_price: { amount: 500, currency: 'jpy' } }
      ],
      confirmed_at: Time.current
    )
  end
 
  before do
    stock = InventoryContext::Stock.new(
      id: 'stock-1',
      product_id: 'product-1',
      total_quantity: 10
    )
    stock_repository.save(stock)
  end
 
  it '在庫を引き当てる' do
    handler.call(event)
 
    stock = stock_repository.find_by_product('product-1')
    expect(stock.available_quantity).to eq(8)  # 10 - 2
  end
end

リナの気づき

Order#confirm! がスッキリした」

以前は確認メソッドの中に在庫処理・メール・ポイント・分析・Slack通知まで詰め込まれていた。ドメインイベントを導入してから、confirm! は「注文を確定する」という純粋なビジネスロジックだけになった。

後続処理は疎結合なハンドラーが担い、新しい処理の追加も Order クラスを触らずにできるようになった。プッシュ通知を追加する際も、既存コードへの変更はゼロ。PushNotificationHandler を作って EventBus.subscribe するだけで完成した。

「これが Open-Closed 原則か。拡張には開いていて、修正には閉じている」

まとめ

  • ドメインイベント = ドメインで起きた出来事を過去形で表現するオブジェクト
  • 集約がイベントを積む = confirm! でイベントをキューに追加、保存後に発行
  • 発行のタイミング = DB保存後(Outboxパターンで原子性を保証)
  • Rails実装 = ActiveSupport::Notificationsで同期処理、ハンドラーで後続処理を分離
  • AWS連携 = EventBridgeでマイクロサービス間の非同期処理
  • テスト = ドメインオブジェクトのイベント発行と、ハンドラーの動作を独立してテスト

次の章では、ユースケースを実装する「アプリケーションサービス」を学ぶ。300行のコントローラーを20行に変える方法を見ていこう。