mybook

イベント駆動アーキテクチャ — 非同期の力

「注文確定したら全部やれ」問題

「ユウキ、注文確定後の処理、今いくつある?」

アヤカの問いに、ユウキはコードを検索した。

# services/order_creation_service.rb の末尾
def finalize_order(order)
  OrderMailer.confirmation(order).deliver_now       # メール送信
  PointService.new(@user).grant_for_order(order)   # ポイント付与
  SlackNotifier.notify("#orders", "注文 #{order.id}")  # Slack通知
  InventoryService.new.update_for_order(order)     # 在庫更新
  RecommendationEngine.new.record_purchase(order)  # レコメンドAI更新
end

「5つ……あ、来週レシートPDF生成も追加予定です」

「そうやって処理が増えていく。今は全部 deliver_now で同期実行だから、メール送信サーバーが遅いとユーザーがずっと待つ」

「確かに……」

「しかも、レコメンドAIが落ちたら注文確定まで失敗する。本来、在庫更新とAIは無関係なのに」

「それに新機能を追加するたびにこのメソッドを変えないといけない。注文サービスが注文に関係ないレコメンドAIの知識を持っている状態はおかしい」

アヤカがうなずいた。「その観察が正しい。これがイベント駆動アーキテクチャを学ぶきっかけだよ」

イベント駆動とは

イベント駆動アーキテクチャでは、「何かが起きた」というイベントを発行し、それを受け取るサブスクライバーが個別に処理する。

Loading diagram...

INFO

イベント駆動の核心は「発行側が受信側を知らない」こと。新しい処理(例:PDF生成)を追加しても、注文サービスのコードは変えない。サブスクライバーを追加するだけ。これを疎結合という。

「発行側と受信側が独立している」ということは:

  • 受信側(メール送信処理)が落ちていても、注文は確定できる
  • 新しい処理を追加するとき、発行側のコードは変わらない
  • 各処理を個別にスケールできる(メール送信だけ増やすなど)

ActiveJob でRails内イベント駆動

まず、Railsの ActiveJob で非同期処理を実装する。

# app/jobs/order_confirmation_job.rb
class OrderConfirmationJob < ApplicationJob
  queue_as :default
  retry_on Net::TimeoutError, wait: :polynomially_longer, attempts: 5
  retry_on ActionMailer::PreviewInterceptorError, attempts: 3
  discard_on ActiveJob::DeserializationError  # レコードが消えた場合は諦める
 
  def perform(order_id)
    order = Order.includes(:user, :order_items).find(order_id)
    OrderMailer.confirmation(order).deliver_now
    Rails.logger.info("確認メール送信完了: order##{order_id}")
  end
end
 
# app/jobs/point_grant_job.rb
class PointGrantJob < ApplicationJob
  queue_as :low_priority
  retry_on StandardError, wait: 30.seconds, attempts: 3
 
  def perform(order_id)
    order = Order.find(order_id)
    PointService.new(order.user).grant_for_order(order)
  end
end
 
# app/jobs/inventory_update_job.rb
class InventoryUpdateJob < ApplicationJob
  queue_as :critical
  retry_on StandardError, attempts: 5, wait: :exponentially_longer
 
  def perform(order_id)
    order = Order.includes(:order_items, :products).find(order_id)
    InventoryService.new.update_for_order(order)
  end
end
 
# app/jobs/recommendation_update_job.rb
class RecommendationUpdateJob < ApplicationJob
  queue_as :low_priority
  discard_on StandardError  # AI更新は失敗しても続行
 
  def perform(order_id)
    order = Order.includes(:order_items).find(order_id)
    RecommendationEngine.new.record_purchase(order)
  end
end
# app/services/order_creation_service.rb
def finalize_order(order)
  # 非同期ジョブをエンキューするだけ
  # 注文サービスは各処理の詳細を知らない
  OrderConfirmationJob.perform_later(order.id)
  PointGrantJob.perform_later(order.id)
  InventoryUpdateJob.perform_later(order.id)
  RecommendationUpdateJob.perform_later(order.id)
  # 注文確定の応答は即座に返る
end

ドメインイベントパターン

より洗練された実装として、ドメインイベントを定義する。

# app/events/base_event.rb
class BaseEvent
  attr_reader :occurred_at, :event_id
 
  def initialize
    @occurred_at = Time.current
    @event_id = SecureRandom.uuid
  end
 
  def self.event_type
    name.underscore.gsub("/", ".")
  end
 
  def to_h
    raise NotImplementedError
  end
end
# app/events/order_completed_event.rb
class OrderCompletedEvent < BaseEvent
  attr_reader :order_id, :user_id, :total_amount, :item_count
 
  def initialize(order)
    super()
    @order_id = order.id
    @user_id = order.user_id
    @total_amount = order.total_amount
    @item_count = order.total_items_count
  end
 
  def to_h
    {
      event_id: @event_id,
      event_type: self.class.event_type,
      order_id: @order_id,
      user_id: @user_id,
      total_amount: @total_amount,
      item_count: @item_count,
      occurred_at: @occurred_at.iso8601
    }
  end
end
 
# app/events/order_cancelled_event.rb
class OrderCancelledEvent < BaseEvent
  attr_reader :order_id, :user_id, :reason
 
  def initialize(order, reason: nil)
    super()
    @order_id = order.id
    @user_id = order.user_id
    @reason = reason
  end
 
  def to_h
    {
      event_id: @event_id,
      event_type: self.class.event_type,
      order_id: @order_id,
      user_id: @user_id,
      reason: @reason,
      occurred_at: @occurred_at.iso8601
    }
  end
end
# app/services/event_bus.rb
class EventBus
  HANDLERS = {
    "order_completed_event" => [
      Handlers::OrderConfirmationHandler,
      Handlers::PointGrantHandler,
      Handlers::InventoryUpdateHandler,
      Handlers::RecommendationUpdateHandler
    ],
    "order_cancelled_event" => [
      Handlers::RefundHandler,
      Handlers::InventoryRestoreHandler,
      Handlers::CancellationNotificationHandler
    ]
  }.freeze
 
  def self.publish(event)
    handlers = HANDLERS[event.class.name.underscore] || []
 
    Rails.logger.info("EventBus: publishing #{event.class.event_type} (#{handlers.size} handlers)")
 
    handlers.each do |handler_class|
      handler_class.new.handle_async(event)
    end
 
    # イベントの記録(監査ログ)
    EventLog.create!(
      event_type: event.class.event_type,
      event_id: event.event_id,
      payload: event.to_h.to_json,
      occurred_at: event.occurred_at
    )
  end
end
# app/handlers/order_confirmation_handler.rb
module Handlers
  class OrderConfirmationHandler
    def handle_async(event)
      OrderConfirmationJob.perform_later(event.order_id)
    end
  end
 
  class PointGrantHandler
    def handle_async(event)
      PointGrantJob.perform_later(event.order_id, event.total_amount)
    end
  end
 
  class InventoryUpdateHandler
    def handle_async(event)
      InventoryUpdateJob.perform_later(event.order_id)
    end
  end
 
  class RecommendationUpdateHandler
    def handle_async(event)
      RecommendationUpdateJob.perform_later(event.order_id)
    end
  end
end
# 使い方:注文サービスはEventBusだけを知っている
class OrderCreationService
  def finalize_order(order)
    EventBus.publish(OrderCompletedEvent.new(order))
    # これだけ!新しい処理はHandlerとEventBusへの登録だけ
  end
end

AWS SQS/SNS での本格的なイベント駆動

Loading diagram...
# app/services/sns_event_publisher.rb
class SnsEventPublisher
  def initialize
    @client = Aws::SNS::Client.new(region: ENV["AWS_REGION"])
    @topic_arn = ENV["ORDER_EVENTS_SNS_TOPIC_ARN"]
  end
 
  def publish(event)
    payload = event.to_h
 
    @client.publish(
      topic_arn: @topic_arn,
      message: payload.to_json,
      subject: event.class.event_type,
      message_attributes: {
        "event_type" => {
          data_type: "String",
          string_value: event.class.event_type
        },
        "event_id" => {
          data_type: "String",
          string_value: event.event_id
        }
      }
    )
 
    Rails.logger.info("SNS published: #{event.class.event_type} (#{event.event_id})")
  rescue Aws::SNS::Errors::ServiceError => e
    Rails.logger.error("SNS publish failed: #{e.message}")
    raise
  end
end
# CloudFormation: SNS → SQS のFan-outアーキテクチャ
Resources:
  OrderEventsTopic:
    Type: AWS::SNS::Topic
    Properties:
      TopicName: order-events
 
  MailQueue:
    Type: AWS::SQS::Queue
    Properties:
      QueueName: order-mail-queue
      VisibilityTimeout: 300
      MessageRetentionPeriod: 86400  # 1日
      RedrivePolicy:
        deadLetterTargetArn: !GetAtt MailDeadLetterQueue.Arn
        maxReceiveCount: 3
 
  MailDeadLetterQueue:
    Type: AWS::SQS::Queue
    Properties:
      QueueName: order-mail-dlq
      MessageRetentionPeriod: 1209600  # 14日
 
  MailQueueSubscription:
    Type: AWS::SNS::Subscription
    Properties:
      TopicArn: !Ref OrderEventsTopic
      Protocol: sqs
      Endpoint: !GetAtt MailQueue.Arn
      FilterPolicy:
        event_type:
          - "order_completed_event"
          - "order_cancelled_event"

SQS メッセージを Rails で処理

# app/jobs/sqs_consumer_job.rb
class SqsConsumerJob < ApplicationJob
  QUEUE_CONFIG = {
    mail: { url: ENV["MAIL_QUEUE_URL"], handler: MailEventHandler },
    point: { url: ENV["POINT_QUEUE_URL"], handler: PointEventHandler },
    inventory: { url: ENV["INVENTORY_QUEUE_URL"], handler: InventoryEventHandler }
  }.freeze
 
  def perform(queue_name)
    config = QUEUE_CONFIG[queue_name.to_sym]
    raise "Unknown queue: #{queue_name}" unless config
 
    client = Aws::SQS::Client.new
    handler = config[:handler].new
 
    loop do
      response = client.receive_message(
        queue_url: config[:url],
        max_number_of_messages: 10,
        wait_time_seconds: 20,         # Long polling
        attribute_names: ["ApproximateReceiveCount"]
      )
 
      response.messages.each do |message|
        process_with_retry(client, message, handler, config[:url])
      end
    end
  end
 
  private
 
  def process_with_retry(client, message, handler, queue_url)
    payload = JSON.parse(message.body)
    event_data = JSON.parse(payload["Message"])
 
    handler.handle(event_data)
 
    client.delete_message(
      queue_url: queue_url,
      receipt_handle: message.receipt_handle
    )
  rescue JSON::ParserError => e
    Rails.logger.error("SQS メッセージのパース失敗: #{e.message}")
    client.delete_message(queue_url: queue_url, receipt_handle: message.receipt_handle)
  rescue => e
    receive_count = message.attributes["ApproximateReceiveCount"].to_i
    Rails.logger.error("SQS処理エラー (試行#{receive_count}回目): #{e.message}")
    raise  # SQSに返却してリトライ
  end
end
# app/handlers/mail_event_handler.rb
class MailEventHandler
  def handle(event_data)
    case event_data["event_type"]
    when "order_completed_event"
      order = Order.find(event_data["order_id"])
      OrderMailer.confirmation(order).deliver_now
    when "order_cancelled_event"
      order = Order.find(event_data["order_id"])
      OrderMailer.cancellation(order, reason: event_data["reason"]).deliver_now
    else
      Rails.logger.warn("不明なイベントタイプ: #{event_data["event_type"]}")
    end
  end
end

デッドレターキューとリトライ戦略

Loading diagram...
# DLQのメッセージを再処理する(手動リカバリ用)
class DlqReprocessJob < ApplicationJob
  def perform(source_queue_url:, target_queue_url:, max_messages: 100)
    dlq_client = Aws::SQS::Client.new
    processed = 0
 
    while processed < max_messages
      response = dlq_client.receive_message(
        queue_url: source_queue_url,
        max_number_of_messages: 10
      )
 
      break if response.messages.empty?
 
      response.messages.each do |msg|
        # 元のキューに再送
        dlq_client.send_message(
          queue_url: target_queue_url,
          message_body: msg.body
        )
 
        dlq_client.delete_message(
          queue_url: source_queue_url,
          receipt_handle: msg.receipt_handle
        )
 
        processed += 1
      end
    end
 
    Rails.logger.info("DLQ再処理完了: #{processed}件")
  end
end

べき等性の実装

# app/jobs/order_confirmation_job.rb
class OrderConfirmationJob < ApplicationJob
  def perform(order_id)
    # べき等キー:同じ注文への確認メールは1回だけ
    idempotency_key = "mail:order_confirmation:#{order_id}"
    return if already_processed?(idempotency_key)
 
    order = Order.find(order_id)
    OrderMailer.confirmation(order).deliver_now
 
    mark_as_processed(idempotency_key)
 
    Rails.logger.info("確認メール送信完了: order##{order_id}")
  end
 
  private
 
  def already_processed?(key)
    Rails.cache.exist?(key)
  end
 
  def mark_as_processed(key)
    Rails.cache.write(key, true, expires_in: 24.hours)
  end
end
# app/jobs/point_grant_job.rb
class PointGrantJob < ApplicationJob
  def perform(order_id, total_amount)
    idempotency_key = "points:grant:#{order_id}"
 
    # Redisのべき等性チェック(SET NX)
    result = Redis.current.set(
      idempotency_key,
      Time.current.iso8601,
      ex: 24 * 60 * 60,  # 24時間
      nx: true            # 存在しない場合のみセット
    )
 
    return Rails.logger.info("ポイント付与スキップ(既処理): order##{order_id}") unless result
 
    order = Order.find(order_id)
    PointService.new(order.user).grant_for_order(order)
  end
end

WARNING

非同期処理になると「注文は完了したのにメールが届かない」という状況が起きうる。デッドレターキューと監視を必ず設定する。べき等性(同じメッセージを2回処理しても結果が同じ)も実装が必要。SQSのAt-Least-Once配信では重複処理が起こりうる。

イベントソーシングへの展開

イベント駆動をさらに発展させると、イベントソーシングになる。状態の代わりにイベントの履歴を保存するパターン。

# app/models/event_log.rb
class EventLog < ApplicationRecord
  validates :event_type, presence: true
  validates :event_id, presence: true, uniqueness: true
  validates :payload, presence: true
 
  scope :for_order, ->(order_id) {
    where("payload->>'order_id' = ?", order_id.to_s)
  }
 
  def event_data
    JSON.parse(payload, symbolize_names: true)
  end
end
# 注文の状態をイベント履歴から再構築
class OrderProjection
  def self.for(order_id)
    events = EventLog.for_order(order_id).order(:occurred_at)
 
    events.reduce({}) do |state, event|
      apply_event(state, event.event_data)
    end
  end
 
  def self.apply_event(state, event)
    case event[:event_type]
    when "order_created_event"
      state.merge(
        status: "pending",
        total_amount: event[:total_amount],
        created_at: event[:occurred_at]
      )
    when "order_completed_event"
      state.merge(status: "confirmed")
    when "order_cancelled_event"
      state.merge(status: "cancelled", cancel_reason: event[:reason])
    else
      state
    end
  end
end

まとめ

処理方式特徴向いている処理
同期(deliver_now)即時実行・失敗がすぐわかるユーザーが結果を待つ処理
非同期(deliver_later)レスポンス速い・再試行可能バックグラウンド処理
イベント駆動(SNS/SQS)疎結合・スケーラブルマイクロサービス間通信

「コードが変わりましたね」ユウキが言った。「注文サービスは注文だけ考えればいい」

「それが目標。各サービスが自分の仕事だけに集中できる。それを単一責任の原則という。イベント駆動はその原則を、サービス間に広げたもの」

「新機能を追加するとき、注文サービスのコードを触らなくていいんですね」

「そう。それが開放閉鎖の原則。拡張に対して開き、変更に対して閉じている。パターンを学ぶと、こういったSOLIDの原則が具体的なコードとして見えてくる」


次章では、データの変換をパイプ&フィルタパターンで整理します。Rackミドルウェアの仕組みと、AWS Step Functionsでの実装を学びます。