ドメインイベント — 変化を通知する
「注文確定したら、在庫を引き当てて、メールを送って、ポイントを付与して、分析データを送って、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
endOrder が注文管理以外の責務を7つも持っている。新しい後続処理が追加されるたびに Order クラスを変更しなければならない。これはOpen-Closed原則(拡張に開いている、修正に閉じている)の違反だ。
ドメインイベントとは
ドメインイベント(Domain Event)は、ドメイン内で起きた「何らかの出来事」を表すオブジェクトだ。
過去形で命名する——「何かが起きた」という事実を表すからだ。
OrderConfirmed(注文が確定した)StockReserved(在庫が引き当てられた)PaymentProcessed(決済が処理された)OrderCancelled(注文がキャンセルされた)
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
endWARNING
イベント発行のタイミングは重要だ。データベース保存の後にイベントを発行しないと、「イベントは発行されたが保存が失敗した」という矛盾状態が発生する。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
endAWS EventBridgeによる非同期処理
FreshCartが成長してマイクロサービス化する場合。
# 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
endINFO
イベントソーシングは強力だが複雑さも高い。「なぜこの状態になったか」の完全な履歴を持てるため、監査要件の厳しい金融系に特に有効だ。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行に変える方法を見ていこう。