イベント駆動アーキテクチャ — 疎結合の実現
「通知サービスを追加したいんですけど、注文サービスのコードを変えずにできますか?」
アオイの質問に、ミサキは微笑んだ。「それがイベント駆動の真骨頂。注文サービスは何も変えなくていい」
イベント駆動 vs 直接呼び出し
Loading diagram...
直接呼び出しの問題:
# 悪い例: 注文サービスが多くのサービスを知っている
class Order < ApplicationRecord
after_commit :notify_all_services, on: :create
private
def notify_all_services
# 通知サービスが追加されるたびにここを修正が必要
NotificationService.send_confirmation(self)
InventoryService.reserve_stock(self)
AnalyticsService.track_order(self)
RecommendationService.update_history(self)
# ↑ 新しいサービスが増えるたびに注文サービスを変更
end
endイベント駆動の解決策:
# 良い例: 注文サービスはイベントを発行するだけ
class Order < ApplicationRecord
after_commit :publish_order_created, on: :create
private
def publish_order_created
# 誰が聞いているかを知らない。発行するだけ
DomainEventPublisher.publish('order.created', self)
# → 通知サービス、在庫サービス、分析サービスは
# それぞれ独立してこのイベントを購読する
end
endドメインイベントの設計
良いドメインイベントを設計するための原則。
# app/events/domain_event.rb
class DomainEvent
attr_reader :event_id, :event_type, :aggregate_id, :aggregate_type,
:occurred_at, :payload, :metadata
def initialize(event_type:, aggregate_id:, aggregate_type:, payload:, metadata: {})
@event_id = SecureRandom.uuid
@event_type = event_type
@aggregate_id = aggregate_id
@aggregate_type = aggregate_type
@occurred_at = Time.current.iso8601(3)
@payload = payload
@metadata = metadata
end
def to_h
{
event_id: event_id,
event_type: event_type,
aggregate_id: aggregate_id,
aggregate_type: aggregate_type,
occurred_at: occurred_at,
payload: payload,
metadata: metadata
}
end
end
# app/events/order_events.rb
module OrderEvents
class OrderCreated < DomainEvent
def initialize(order:, user_agent: nil)
super(
event_type: 'order.created',
aggregate_id: order.id,
aggregate_type: 'Order',
payload: {
order_id: order.id,
user_id: order.user_id,
total_amount_cents: order.total_amount_cents,
currency: 'JPY',
items: order.order_items.map { |item|
{
product_id: item.product_id,
product_name: item.product_name_snapshot,
quantity: item.quantity,
unit_price_cents: item.unit_price_snapshot_cents
}
},
shipping_address: order.shipping_address
},
metadata: {
user_agent: user_agent,
service: 'order-service',
version: '1.0'
}
)
end
end
endAWS EventBridge: イベントバスの中心
EventBridge は AWS のマネージドイベントバス。ルールベースでイベントを各サービスに振り分ける。
# app/services/event_bridge_publisher.rb
class EventBridgePublisher
CLIENT = Aws::EventBridge::Client.new(region: 'ap-northeast-1')
EVENT_BUS_NAME = ENV.fetch('EVENT_BUS_NAME', 'shopnova-events')
def self.publish(domain_event)
CLIENT.put_events(
entries: [{
event_bus_name: EVENT_BUS_NAME,
source: "shopnova.#{domain_event.aggregate_type.downcase}",
detail_type: domain_event.event_type,
detail: domain_event.to_h.to_json,
time: Time.parse(domain_event.occurred_at)
}]
)
end
def self.publish_batch(domain_events)
# EventBridgeは1回のAPIコールで最大10イベント
domain_events.each_slice(10) do |batch|
CLIENT.put_events(
entries: batch.map { |event|
{
event_bus_name: EVENT_BUS_NAME,
source: "shopnova.#{event.aggregate_type.downcase}",
detail_type: event.event_type,
detail: event.to_h.to_json
}
}
)
end
end
end# EventBridge ルール定義(CloudFormation)
NotificationServiceRule:
Type: AWS::Events::Rule
Properties:
EventBusName: shopnova-events
EventPattern:
source:
- shopnova.order
detail-type:
- order.created
- order.cancelled
- order.shipped
Targets:
- Id: NotificationSQSQueue
Arn: !GetAtt NotificationQueue.Arn
InventoryServiceRule:
Type: AWS::Events::Rule
Properties:
EventBusName: shopnova-events
EventPattern:
source:
- shopnova.order
detail-type:
- order.created
- order.cancelled
Targets:
- Id: InventorySQSQueue
Arn: !GetAtt InventoryQueue.Arn
AnalyticsServiceRule:
Type: AWS::Events::Rule
Properties:
EventBusName: shopnova-events
EventPattern:
source:
- prefix: shopnova. # 全サービスのイベントを受信
Targets:
- Id: AnalyticsFirehose
Arn: !GetAtt AnalyticsFirehose.Arn # Kinesis Firehoseに送ってS3に蓄積イベントソーシング入門
イベントソーシングは、状態の変化を「イベントの列」として保存するパターン。
# 通常の設計: 現在の状態だけを保存
class Order < ApplicationRecord
# status: 'pending' → 'paid' → 'shipped' → 'delivered'
# 過去の状態は見えない
end
# イベントソーシング: イベントの列を保存
class OrderEvent < ApplicationRecord
# order_id, event_type, payload, occurred_at
# event_type: 'order_placed', 'payment_received', 'shipped', 'delivered'
# 現在の状態 = イベントを再生して導出
end
# イベントから現在の状態を再構築
class OrderProjection
def self.build(order_id)
events = OrderEvent.where(order_id: order_id).order(:occurred_at)
events.reduce({}) do |state, event|
apply_event(state, event)
end
end
def self.apply_event(state, event)
case event.event_type
when 'order_placed'
state.merge(
id: event.order_id,
status: 'pending',
items: event.payload['items'],
created_at: event.occurred_at
)
when 'payment_received'
state.merge(status: 'paid', paid_at: event.occurred_at)
when 'shipped'
state.merge(
status: 'shipped',
tracking_number: event.payload['tracking_number'],
shipped_at: event.occurred_at
)
else
state
end
end
endINFO
イベントソーシングを使うと「2024年12月1日時点の注文の状態は?」という時点の状態を復元できる。監査ログ・デバッグ・タイムトラベルデバッグに強い。ただし複雑性も増すため、必要性をよく検討すること。
通知サービスの実装
EventBridgeからのイベントを受けて、適切な通知を送信する。
# app/workers/notification_worker.rb
class NotificationWorker
include Shoryuken::Worker
shoryuken_options queue: ENV.fetch('NOTIFICATION_QUEUE_URL')
HANDLERS = {
'order.created' => :handle_order_created,
'order.shipped' => :handle_order_shipped,
'order.cancelled' => :handle_order_cancelled,
'product.back_in_stock' => :handle_back_in_stock
}.freeze
def perform(sqs_msg, body)
# EventBridgeからのメッセージはSNS経由でラップされている
sns_message = JSON.parse(body)
event = JSON.parse(sns_message['Message'])
handler_method = HANDLERS[event['detail-type']]
if handler_method
send(handler_method, event['detail'])
else
Rails.logger.info("Unhandled event type: #{event['detail-type']}")
end
sqs_msg.delete
end
private
def handle_order_created(detail)
user = UserServiceClient.find(detail['payload']['user_id'])
# メール通知
NotificationMailer.order_confirmation(
to: user['email'],
order_id: detail['payload']['order_id'],
items: detail['payload']['items'],
total: detail['payload']['total_amount_cents'] / 100.0
).deliver_later
# プッシュ通知(FCM)
if user['fcm_token']
PushNotificationService.send(
token: user['fcm_token'],
title: 'ご注文を承りました',
body: "注文番号: #{detail['payload']['order_id']}"
)
end
end
def handle_order_shipped(detail)
user = UserServiceClient.find(detail['payload']['user_id'])
NotificationMailer.shipment_notification(
to: user['email'],
tracking_number: detail['payload']['tracking_number']
).deliver_later
end
endイベントのスキーマ管理
イベントのスキーマが変更されたとき、購読者が壊れないようにする。
# イベントスキーマのバージョニング
class OrderEvents
# v1: シンプルな構造
class OrderCreatedV1 < DomainEvent
def initialize(order:)
super(
event_type: 'order.created.v1',
payload: { order_id: order.id, total: order.total }
)
end
end
# v2: 詳細な構造(後方互換性あり)
class OrderCreatedV2 < DomainEvent
def initialize(order:)
super(
event_type: 'order.created.v2',
payload: {
order_id: order.id,
total_amount_cents: order.total_amount_cents,
currency: 'JPY',
items: order.items.map(&:to_event_payload)
# v1のフィールドも残す(後方互換性)
}
)
end
end
end# EventBridgeスキーマレジストリの活用
# AWS EventBridgeのSchema Registryでスキーマを管理
schema:
name: shopnova-order-created
version: "2"
content: |
{
"$schema": "http://json-schema.org/draft-04/schema#",
"type": "object",
"properties": {
"order_id": { "type": "string" },
"user_id": { "type": "string" },
"total_amount_cents": { "type": "integer" },
"currency": { "type": "string" },
"items": {
"type": "array",
"items": {
"type": "object",
"properties": {
"product_id": { "type": "string" },
"quantity": { "type": "integer" }
}
}
}
},
"required": ["order_id", "user_id", "total_amount_cents", "items"]
}デッドレターキューとエラーハンドリング
# 全てのSQSキューにDLQを設定
OrderEventsSQS:
Type: AWS::SQS::Queue
Properties:
QueueName: shopnova-order-events
VisibilityTimeout: 60
RedrivePolicy:
deadLetterTargetArn: !GetAtt OrderEventsDLQ.Arn
maxReceiveCount: 3 # 3回失敗でDLQに移動
OrderEventsDLQ:
Type: AWS::SQS::Queue
Properties:
QueueName: shopnova-order-events-dlq
MessageRetentionPeriod: 1209600 # 14日間保持# DLQを監視してアラートを送信
class DlqMonitorJob < ApplicationJob
queue_as :monitoring
def perform
sqs = Aws::SQS::Client.new
dlq_url = ENV.fetch('ORDER_EVENTS_DLQ_URL')
attributes = sqs.get_queue_attributes(
queue_url: dlq_url,
attribute_names: ['ApproximateNumberOfMessages']
)
message_count = attributes.attributes['ApproximateNumberOfMessages'].to_i
if message_count > 0
AlertService.notify(
severity: 'warning',
message: "DLQに #{message_count} 件のメッセージが溜まっています",
queue: dlq_url
)
end
end
endまとめ
「注文サービスを1行も変えずに通知サービスを追加できました」とアオイが報告した。
ミサキはうなずいた。「これがイベント駆動の力。サービスが増えても、既存のコードは変わらない」
イベント駆動の原則:
1. 発行者は購読者を知らない(疎結合)
2. イベントは過去形で命名する(order.created ✓、create_order ✗)
3. イベントにはビジネスの意味を込める
4. スキーマは後方互換を保ちながら進化させる
5. DLQで失敗イベントを捕捉する
WARNING
イベント駆動で気をつけること: イベントの順序保証は難しい。order.created より order.cancelled が先に届くことがある。各サービスは順序が前後することを前提に実装する。
次章では、サービスが増えた環境で「どのサービスがどこにいるか」を見つけるサービスディスカバリを学ぶ。