サーガパターン — 分散トランザクションを管理する
「マイクロサービスでトランザクションが使えない」
「アヤカさん、大変なことに気づきました」
ユウキが青ざめた顔でアヤカのデスクに来た。
「注文サービスと在庫サービスが別のデータベースなんですよ。注文を作ったとき、在庫を減らす前にサーバーがクラッシュしたら……?」
「在庫が減らないまま、注文だけ作られる」
「ですよね! BEGIN TRANSACTION って別DBにまたがれないし……」
アヤカがうなずいた。「それがマイクロサービスの最大の難しさ。分散トランザクション。サーガパターンで解決する」
「サーガって何ですか?」
「北欧の長編叙事詩の意味。長い旅の物語。分散トランザクションも、複数のステップを経る長い旅だから」
サーガパターンとは
サーガは、複数サービスにまたがる長時間トランザクションを、小さなローカルトランザクションに分解するパターン。各ステップが失敗したとき、補償トランザクション(逆の操作)で一貫性を回復する。
INFO
サーガは「完全に成功するか、補償処理で元に戻す」ことを保証する。ただし補償トランザクションが実行中は、他のサービスから見ると不整合な状態が一時的に存在する(結果整合性)。ACIDトランザクションと異なり、BASE(Basically Available, Soft state, Eventually consistent)の性質を持つ。
オーケストレーション型サーガ
中央の指揮者(オーケストレーター)が各ステップを管理する。
# app/sagas/order_fulfillment_saga.rb
class OrderFulfillmentSaga
STEPS = [
:reserve_inventory,
:process_payment,
:schedule_delivery
].freeze
def initialize(order_id)
@order_id = order_id
@saga_state = SagaState.find_or_create_by!(
order_id: order_id,
saga_type: self.class.name
)
end
def execute
validate_initial_state!
STEPS.each do |step|
next if @saga_state.completed_steps.include?(step.to_s)
Rails.logger.info("Saga[#{@order_id}]: executing #{step}")
result = send(step)
if result.success?
@saga_state.mark_step_completed!(step)
Rails.logger.info("Saga[#{@order_id}]: #{step} completed")
else
Rails.logger.error("Saga[#{@order_id}]: #{step} failed: #{result.error}")
compensate_from(step)
raise SagaFailedError, "#{step} 失敗: #{result.error}"
end
end
@saga_state.mark_completed!
Rails.logger.info("Saga[#{@order_id}]: completed successfully")
end
private
def validate_initial_state!
raise SagaAlreadyCompletedError if @saga_state.completed?
raise SagaFailedError if @saga_state.failed?
end
def reserve_inventory
items = order.order_items.map { |i|
{ product_id: i.product_id, qty: i.quantity }
}
response = InventoryServiceClient.new.reserve(
order_id: @order_id,
items: items,
idempotency_key: "inventory:#{@order_id}"
)
ServiceResult.new(success: response.success?, error: response.error_message,
data: { reservation_id: response.reservation_id })
end
def process_payment
payment_method = order.user.default_payment_method
response = PaymentServiceClient.new.charge(
order_id: @order_id,
amount: order.total_amount,
user_id: order.user_id,
payment_method_id: payment_method.id,
idempotency_key: "payment:#{@order_id}"
)
ServiceResult.new(success: response.success?, error: response.error_message,
data: { payment_id: response.payment_id })
end
def schedule_delivery
response = DeliveryServiceClient.new.schedule(
order_id: @order_id,
address: order.shipping_address,
items: order.order_items.map { |i|
{ product_id: i.product_id, qty: i.quantity, weight: i.product.weight }
},
idempotency_key: "delivery:#{@order_id}"
)
ServiceResult.new(success: response.success?, error: response.error_message,
data: { delivery_id: response.delivery_id })
end
def compensate_from(failed_step)
completed = @saga_state.completed_steps.map(&:to_sym)
step_index = STEPS.index(failed_step)
# 完了済みのステップを逆順に補償
completed.reverse.each do |step|
next if STEPS.index(step) >= step_index
begin
send("compensate_#{step}")
Rails.logger.info("Saga[#{@order_id}]: compensated #{step}")
rescue => e
Rails.logger.error("Saga[#{@order_id}]: compensation failed for #{step}: #{e.message}")
# 補償の失敗は人間の介入が必要
PagerDutyAlertJob.perform_later(
"Saga compensation failed",
order_id: @order_id,
step: step,
error: e.message
)
end
end
@saga_state.mark_failed!(failed_step)
end
def compensate_reserve_inventory
InventoryServiceClient.new.release(
order_id: @order_id,
idempotency_key: "inventory_release:#{@order_id}"
)
end
def compensate_process_payment
PaymentServiceClient.new.refund(
order_id: @order_id,
idempotency_key: "refund:#{@order_id}"
)
end
def order
@order ||= Order.includes(:user, :order_items, :products).find(@order_id)
end
end# app/models/saga_state.rb
class SagaState < ApplicationRecord
serialize :completed_steps, type: Array, coder: JSON
validates :order_id, presence: true
validates :saga_type, presence: true
validates :status, inclusion: { in: %w[in_progress completed failed] }
scope :in_progress, -> { where(status: "in_progress") }
scope :failed, -> { where(status: "failed") }
scope :stale, -> { in_progress.where("updated_at < ?", 30.minutes.ago) }
def mark_step_completed!(step)
with_lock do
completed_steps << step.to_s
save!
end
end
def mark_completed!
update!(status: "completed", completed_at: Time.current)
end
def mark_failed!(step)
update!(
status: "failed",
failed_at_step: step.to_s,
failed_at: Time.current
)
end
def completed?
status == "completed"
end
def failed?
status == "failed"
end
end# マイグレーション
class CreateSagaStates < ActiveRecord::Migration[8.0]
def change
create_table :saga_states do |t|
t.bigint :order_id, null: false
t.string :saga_type, null: false
t.string :status, null: false, default: "in_progress"
t.text :completed_steps, default: "[]"
t.string :failed_at_step
t.json :context_data # ステップ間でデータを受け渡す
t.timestamp :completed_at
t.timestamp :failed_at
t.timestamps
end
add_index :saga_states, [:order_id, :saga_type], unique: true
add_index :saga_states, :status
add_index :saga_states, :updated_at # staleなサーガの検索用
end
endコレオグラフィ型サーガ
各サービスがイベントを通じて連鎖する。指揮者は不在。
# 在庫サービス:OrderCreated イベントを受信して処理
class InventoryEventHandler
def handle_order_created(event)
order_data = event.data
begin
items = order_data[:items]
reserve_all(items, order_data[:order_id])
EventBus.publish(InventoryReservedEvent.new(
order_id: order_data[:order_id],
reservation_ids: items.map { |i| i[:reservation_id] }
))
rescue InsufficientStockError => e
EventBus.publish(InventoryReservationFailedEvent.new(
order_id: order_data[:order_id],
reason: e.message,
failed_items: e.failed_items
))
end
end
private
def reserve_all(items, order_id)
# 全部成功するか全部失敗するかの実装
reservations = []
items.each do |item|
product = Product.lock.find(item[:product_id])
raise InsufficientStockError.new(product.name, item[:quantity]) unless product.in_stock?(item[:quantity])
product.decrement!(:stock, item[:quantity])
reservations << {
product_id: item[:product_id],
quantity: item[:quantity],
reserved_at: Time.current
}
end
InventoryReservation.create!(
order_id: order_id,
reservations: reservations.to_json
)
end
end# 注文サービス:InventoryReservationFailed イベントを受信して補償
class OrderEventHandler
def handle_inventory_reservation_failed(event)
order = Order.find(event.data[:order_id])
order.cancel!(reason: "在庫不足: #{event.data[:reason]}")
OrderMailer.cancellation(
order,
reason: event.data[:reason]
).deliver_later
end
def handle_payment_failed(event)
order = Order.find(event.data[:order_id])
# 在庫を解放するよう指示
EventBus.publish(InventoryReleaseRequestedEvent.new(
order_id: event.data[:order_id],
reason: "決済失敗による在庫解放"
))
order.cancel!(reason: "決済失敗: #{event.data[:reason]}")
OrderMailer.payment_failed(order).deliver_later
end
endAWS Step Functions でのサーガ
{
"Comment": "注文フルフィルメントサーガ",
"StartAt": "ReserveInventory",
"States": {
"ReserveInventory": {
"Type": "Task",
"Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:reserve-inventory",
"Parameters": {
"order_id.$": "$.order_id",
"items.$": "$.items"
},
"ResultPath": "$.inventory_result",
"Next": "ProcessPayment",
"Catch": [
{
"ErrorEquals": ["InsufficientStockError"],
"ResultPath": "$.error",
"Next": "HandleInventoryError"
}
],
"Retry": [
{
"ErrorEquals": ["Lambda.TooManyRequestsException"],
"IntervalSeconds": 1,
"MaxAttempts": 3,
"BackoffRate": 2
}
]
},
"ProcessPayment": {
"Type": "Task",
"Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:process-payment",
"Parameters": {
"order_id.$": "$.order_id",
"amount.$": "$.total_amount",
"user_id.$": "$.user_id"
},
"ResultPath": "$.payment_result",
"Next": "ScheduleDelivery",
"Catch": [
{
"ErrorEquals": ["PaymentError", "InsufficientFundsError"],
"ResultPath": "$.error",
"Next": "CompensateInventory"
}
]
},
"ScheduleDelivery": {
"Type": "Task",
"Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:schedule-delivery",
"ResultPath": "$.delivery_result",
"Next": "OrderCompleted",
"Catch": [
{
"ErrorEquals": ["DeliveryError"],
"ResultPath": "$.error",
"Next": "CompensatePayment"
}
]
},
"CompensatePayment": {
"Type": "Task",
"Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:refund-payment",
"Parameters": {
"order_id.$": "$.order_id",
"payment_id.$": "$.payment_result.payment_id"
},
"Next": "CompensateInventory"
},
"CompensateInventory": {
"Type": "Task",
"Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:release-inventory",
"Parameters": {
"order_id.$": "$.order_id"
},
"Next": "OrderFailed"
},
"HandleInventoryError": {
"Type": "Task",
"Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:notify-inventory-failed",
"Next": "OrderFailed"
},
"OrderCompleted": {
"Type": "Succeed"
},
"OrderFailed": {
"Type": "Fail",
"Error": "OrderFulfillmentFailed",
"Cause": "注文フルフィルメントに失敗しました"
}
}
}サーガのテスト
# spec/sagas/order_fulfillment_saga_spec.rb
RSpec.describe OrderFulfillmentSaga do
let(:order) { create(:order, :with_items, total_amount: 5000) }
let(:saga) { described_class.new(order.id) }
let(:inventory_client) { instance_double(InventoryServiceClient) }
let(:payment_client) { instance_double(PaymentServiceClient) }
let(:delivery_client) { instance_double(DeliveryServiceClient) }
before do
allow(InventoryServiceClient).to receive(:new).and_return(inventory_client)
allow(PaymentServiceClient).to receive(:new).and_return(payment_client)
allow(DeliveryServiceClient).to receive(:new).and_return(delivery_client)
end
context "全ステップ成功の場合" do
before do
allow(inventory_client).to receive(:reserve)
.and_return(double(success?: true, error_message: nil, reservation_id: "r-123"))
allow(payment_client).to receive(:charge)
.and_return(double(success?: true, error_message: nil, payment_id: "p-456"))
allow(delivery_client).to receive(:schedule)
.and_return(double(success?: true, error_message: nil, delivery_id: "d-789"))
end
it "サーガが完了状態になる" do
saga.execute
expect(SagaState.find_by(order_id: order.id).status).to eq("completed")
end
it "全ステップが完了として記録される" do
saga.execute
state = SagaState.find_by(order_id: order.id)
expect(state.completed_steps).to contain_exactly(
"reserve_inventory", "process_payment", "schedule_delivery"
)
end
end
context "決済が失敗した場合" do
before do
allow(inventory_client).to receive(:reserve)
.and_return(double(success?: true, error_message: nil, reservation_id: "r-123"))
allow(payment_client).to receive(:charge)
.and_return(double(success?: false, error_message: "カードが拒否されました", payment_id: nil))
allow(inventory_client).to receive(:release).and_return(true)
end
it "SagaFailedError を発生させる" do
expect { saga.execute }.to raise_error(SagaFailedError, /payment/)
end
it "在庫が解放される(補償トランザクション)" do
saga.execute rescue nil
expect(inventory_client).to have_received(:release)
end
it "サーガが失敗状態になる" do
saga.execute rescue nil
state = SagaState.find_by(order_id: order.id)
expect(state.status).to eq("failed")
expect(state.failed_at_step).to eq("process_payment")
end
end
context "配送予約が失敗した場合" do
before do
allow(inventory_client).to receive(:reserve)
.and_return(double(success?: true, error_message: nil, reservation_id: "r-123"))
allow(payment_client).to receive(:charge)
.and_return(double(success?: true, error_message: nil, payment_id: "p-456"))
allow(delivery_client).to receive(:schedule)
.and_return(double(success?: false, error_message: "配送エリア外", delivery_id: nil))
allow(payment_client).to receive(:refund).and_return(true)
allow(inventory_client).to receive(:release).and_return(true)
end
it "決済が返金される" do
saga.execute rescue nil
expect(payment_client).to have_received(:refund)
end
it "在庫が解放される" do
saga.execute rescue nil
expect(inventory_client).to have_received(:release)
end
end
endWARNING
補償トランザクション自体が失敗するケースも考慮が必要。補償の補償は存在しない。この場合は人間の介入が必要なケースとして、アラートとDead Letter Queueに記録する設計にする。また、サーガは長時間実行されることがあるため、タイムアウトと再開機能も設計に含めること。
ステール状態のサーガを監視する
# app/jobs/saga_monitor_job.rb
class SagaMonitorJob < ApplicationJob
queue_as :monitoring
def perform
stale_sagas = SagaState.stale # 30分以上 in_progress のまま
stale_sagas.each do |state|
Rails.logger.warn("Stale saga detected: order##{state.order_id}, step: #{state.completed_steps.last}")
# PagerDuty/SNS に通知
SagaAlertJob.perform_later(
order_id: state.order_id,
stuck_at_step: state.completed_steps.last,
duration_minutes: ((Time.current - state.updated_at) / 60).round
)
end
end
end
# config/initializers/sidekiq.rb に定期実行を登録
Sidekiq::Cron::Job.create(
name: "Saga Monitor",
cron: "*/5 * * * *", # 5分ごと
class: "SagaMonitorJob"
)オーケストレーション vs コレオグラフィ
| 観点 | オーケストレーション | コレオグラフィ |
|---|---|---|
| 管理の明確さ | 一箇所で全フロー把握 | 各サービスに分散 |
| 結合度 | 中央オーケストレーターに依存 | サービス同士は疎結合 |
| デバッグのしやすさ | フローが追いやすい | イベントログを追う必要 |
| 障害の特定 | どのステップで失敗したか明確 | 分散トレーシングが必要 |
| スケーラビリティ | オーケストレーターがボトルネックになりうる | 各サービスが独立にスケール |
| 向いている規模 | 中規模・3〜5サービス | 大規模・多数サービス |
「どっちを使えばいいですか?」
「最初はオーケストレーション」アヤカが答えた。「フローが一箇所にあるから、デバッグしやすい。Step Functionsなら実行履歴が全部残る。サービスが10個を超えたら、コレオグラフィを考え始める」
「サーガを使うのはどんなときですか?」
「複数のマイクロサービスをまたがって、一貫性を保ちたいとき。1つのサービス内なら普通のDBトランザクションでいい。2〜3サービスをまたぐ処理が出てきたらサーガを検討する」
次章では、モノリスをマイクロサービスに移行する現実的な方法、ストラングラーフィグパターンを学びます。