mybook

イベント駆動アーキテクチャ — 疎結合の実現

「通知サービスを追加したいんですけど、注文サービスのコードを変えずにできますか?」

アオイの質問に、ミサキは微笑んだ。「それがイベント駆動の真骨頂。注文サービスは何も変えなくていい」


イベント駆動 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
end

AWS 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
end

INFO

イベントソーシングを使うと「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 が先に届くことがある。各サービスは順序が前後することを前提に実装する。

次章では、サービスが増えた環境で「どのサービスがどこにいるか」を見つけるサービスディスカバリを学ぶ。