mybook

サーガパターン — 分散トランザクションを管理する

「マイクロサービスでトランザクションが使えない」

「アヤカさん、大変なことに気づきました」

ユウキが青ざめた顔でアヤカのデスクに来た。

「注文サービスと在庫サービスが別のデータベースなんですよ。注文を作ったとき、在庫を減らす前にサーバーがクラッシュしたら……?」

「在庫が減らないまま、注文だけ作られる」

「ですよね! BEGIN TRANSACTION って別DBにまたがれないし……」

アヤカがうなずいた。「それがマイクロサービスの最大の難しさ。分散トランザクション。サーガパターンで解決する」

「サーガって何ですか?」

「北欧の長編叙事詩の意味。長い旅の物語。分散トランザクションも、複数のステップを経る長い旅だから」

サーガパターンとは

サーガは、複数サービスにまたがる長時間トランザクションを、小さなローカルトランザクションに分解するパターン。各ステップが失敗したとき、補償トランザクション(逆の操作)で一貫性を回復する。

Loading diagram...

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

コレオグラフィ型サーガ

各サービスがイベントを通じて連鎖する。指揮者は不在。

Loading diagram...
# 在庫サービス: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
end

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

WARNING

補償トランザクション自体が失敗するケースも考慮が必要。補償の補償は存在しない。この場合は人間の介入が必要なケースとして、アラートと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サービスをまたぐ処理が出てきたらサーガを検討する」


次章では、モノリスをマイクロサービスに移行する現実的な方法、ストラングラーフィグパターンを学びます。