イベント駆動アーキテクチャ — 非同期の力
「注文確定したら全部やれ」問題
「ユウキ、注文確定後の処理、今いくつある?」
アヤカの問いに、ユウキはコードを検索した。
# 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の知識を持っている状態はおかしい」
アヤカがうなずいた。「その観察が正しい。これがイベント駆動アーキテクチャを学ぶきっかけだよ」
イベント駆動とは
イベント駆動アーキテクチャでは、「何かが起きた」というイベントを発行し、それを受け取るサブスクライバーが個別に処理する。
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
endAWS SQS/SNS での本格的なイベント駆動
# 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デッドレターキューとリトライ戦略
# 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
endWARNING
非同期処理になると「注文は完了したのにメールが届かない」という状況が起きうる。デッドレターキューと監視を必ず設定する。べき等性(同じメッセージを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での実装を学びます。