mybook

サーガパターン — 分散トランザクション

「注文が確定したとき、在庫を減らして、決済を処理して、ポイントを付与する — これって全部成功か全部失敗にしたいですよね」

ミサキがうなずく前に、シニアエンジニアのアオイが言葉を続けた。「でも、それぞれ別サービスだから、BEGIN TRANSACTION が使えない」

「そこがマイクロサービスの一番難しいところ」とミサキは答えた。「サーガパターンで解決する」


分散トランザクションの問題

モノリスなら簡単だった処理が、マイクロサービスでは複雑になる。

# モノリス: 1つのトランザクションで全て完結
ActiveRecord::Base.transaction do
  order = Order.create!(order_params)
  inventory.decrement!(quantity: order.quantity)
  payment.charge!(amount: order.total)
  user.add_points!(order.total / 100)
end
# どこかで失敗 → 全部ロールバック
 
# マイクロサービス: それぞれ別のサービス・DB
# 注文サービス    → 注文作成(orders DB)
# 在庫サービス   → 在庫減算(inventory DB)
# 決済サービス   → 決済処理(payment DB)
# ユーザーサービス → ポイント付与(users DB)
# → 途中で失敗したら? 整合性を取り戻すには?
Loading diagram...

サーガパターン

**サーガ(Saga)**は、複数のサービスにまたがるビジネストランザクションを、**補償トランザクション(Compensating Transaction)**によって整合性を保つパターン。

2種類の実装方式がある:

方式特徴
Choreography(振り付け)サービス同士がイベントで連携。中央コントローラなし
Orchestration(指揮)中央のオーケストレーターがサービスに命令

Choreography サーガ

各サービスがイベントを発行・購読することで流れを制御する。

Loading diagram...
# 注文サービス: 注文を作成してイベントを発行
class CreateOrderSaga
  def call(user_id:, items:, payment_method:)
    order = Order.create!(
      user_id: user_id,
      status: 'pending',
      items_data: items,
      payment_method: payment_method
    )
 
    # 在庫サービスに在庫予約を依頼
    SnsPublisher.publish(
      topic_arn: ENV.fetch('ORDER_EVENTS_TOPIC_ARN'),
      event_type: 'order.created',
      payload: {
        saga_id: order.id,  # saga_id として order_id を使用
        order_id: order.id,
        user_id: user_id,
        items: items
      }
    )
 
    order
  end
end
 
# 在庫サービス: order.created を受けて在庫予約
class InventoryWorker
  include Shoryuken::Worker
  shoryuken_options queue: ENV.fetch('INVENTORY_QUEUE_URL')
 
  def perform(sqs_msg, body)
    event = JSON.parse(body)
    return unless event['event_type'] == 'order.created'
 
    saga_id = event['payload']['saga_id']
    items = event['payload']['items']
 
    begin
      # 在庫を予約(減算ではなく予約状態に)
      reservation = InventoryReservation.create_for_order!(
        saga_id: saga_id,
        items: items
      )
 
      # 成功: 次のステップへ
      SnsPublisher.publish(
        topic_arn: ENV.fetch('INVENTORY_EVENTS_TOPIC_ARN'),
        event_type: 'inventory.reserved',
        payload: {
          saga_id: saga_id,
          reservation_id: reservation.id,
          order_id: event['payload']['order_id']
        }
      )
    rescue InsufficientStockError => e
      # 失敗: サーガを補償(注文をキャンセル)
      SnsPublisher.publish(
        topic_arn: ENV.fetch('INVENTORY_EVENTS_TOPIC_ARN'),
        event_type: 'inventory.reservation_failed',
        payload: {
          saga_id: saga_id,
          order_id: event['payload']['order_id'],
          reason: e.message
        }
      )
    end
 
    sqs_msg.delete
  end
end
 
# 注文サービス: 失敗イベントを受けてロールバック
class OrderFailureWorker
  include Shoryuken::Worker
 
  def perform(sqs_msg, body)
    event = JSON.parse(body)
 
    case event['event_type']
    when 'inventory.reservation_failed', 'payment.failed'
      # 補償トランザクション: 注文をキャンセル
      order = Order.find_by(id: event['payload']['order_id'])
      order&.update!(status: 'cancelled', cancelled_reason: event['payload']['reason'])
 
      NotificationService.send_cancellation_email(order)
    end
 
    sqs_msg.delete
  end
end

Choreography の問題点

メリット:
  ✓ シンプル、中央コントローラが不要
  ✓ 各サービスが疎結合

デメリット:
  ✗ 全体の流れが追いにくい(どのサービスが何を担当?)
  ✗ 循環依存が発生しやすい
  ✗ デバッグが困難(障害の原因を追うのが大変)

Orchestration サーガ

中央のオーケストレーター(AWS Step Functions)が各サービスに指示する。

Loading diagram...
{
  "Comment": "ShopNova Order Processing Saga",
  "StartAt": "ReserveInventory",
  "States": {
    "ReserveInventory": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxxx:function:reserve-inventory",
      "ResultPath": "$.inventoryResult",
      "Catch": [{
        "ErrorEquals": ["InsufficientStockError"],
        "ResultPath": "$.error",
        "Next": "CancelOrder"
      }],
      "Next": "ProcessPayment"
    },
    "ProcessPayment": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxxx:function:process-payment",
      "ResultPath": "$.paymentResult",
      "Catch": [{
        "ErrorEquals": ["PaymentError"],
        "ResultPath": "$.error",
        "Next": "ReleaseInventory"
      }],
      "Next": "AddPoints"
    },
    "AddPoints": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxxx:function:add-points",
      "Next": "CompleteOrder"
    },
    "CompleteOrder": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxxx:function:complete-order",
      "End": true
    },
    "ReleaseInventory": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxxx:function:release-inventory",
      "Next": "CancelOrder"
    },
    "CancelOrder": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxxx:function:cancel-order",
      "End": true
    }
  }
}
# Lambda: 在庫予約
# lambda/reserve_inventory/handler.rb
def handler(event:, context:)
  order_id = event['order_id']
  items = event['items']
 
  begin
    reservation = InventoryService.reserve(order_id: order_id, items: items)
    {
      status: 'success',
      reservation_id: reservation.id
    }
  rescue InsufficientStockError => e
    raise e  # Step Functionsがキャッチして補償フローに移行
  end
end
 
# Lambda: 補償 - 在庫解放
# lambda/release_inventory/handler.rb
def handler(event:, context:)
  order_id = event['order_id']
 
  reservation = InventoryReservation.find_by(order_id: order_id)
  reservation&.release!
 
  { status: 'released' }
end

Rails からStep Functions を起動

# app/services/order_saga_starter.rb
class OrderSagaStarter
  STEP_FUNCTIONS = Aws::States::Client.new(region: 'ap-northeast-1')
  STATE_MACHINE_ARN = ENV.fetch('ORDER_SAGA_STATE_MACHINE_ARN')
 
  def self.start(order)
    input = {
      order_id: order.id,
      user_id: order.user_id,
      items: order.order_items.map { |item|
        {
          product_id: item.product_id,
          quantity: item.quantity,
          price_cents: item.unit_price_snapshot_cents
        }
      },
      total_amount_cents: order.total_amount_cents,
      payment_method: order.payment_method
    }
 
    execution = STEP_FUNCTIONS.start_execution(
      state_machine_arn: STATE_MACHINE_ARN,
      name: "order-#{order.id}-#{Time.current.to_i}",
      input: input.to_json
    )
 
    order.update!(saga_execution_arn: execution.execution_arn)
    execution
  end
 
  def self.status(order)
    execution = STEP_FUNCTIONS.describe_execution(
      execution_arn: order.saga_execution_arn
    )
    execution.status  # RUNNING, SUCCEEDED, FAILED
  end
end
 
# 注文作成APIから呼び出す
class OrdersController < ApplicationController
  def create
    order = Order.create!(order_params)
    OrderSagaStarter.start(order)
    render json: { order_id: order.id, status: 'processing' }, status: :accepted
  end
end

サーガの状態管理

# db/migrate/xxxx_create_saga_executions.rb
class CreateSagaExecutions < ActiveRecord::Migration[7.2]
  def change
    create_table :saga_executions, id: :uuid do |t|
      t.string :saga_type, null: false
      t.string :status, null: false, default: 'running'
      t.jsonb :context, null: false, default: {}
      t.jsonb :completed_steps, null: false, default: []
      t.jsonb :compensation_steps, null: false, default: []
      t.text :failure_reason
      t.timestamps
    end
 
    add_index :saga_executions, :status
  end
end
 
# app/models/saga_execution.rb
class SagaExecution < ApplicationRecord
  STATUSES = %w[running succeeded failed compensating compensated].freeze
 
  validates :status, inclusion: { in: STATUSES }
 
  def complete_step!(step_name, result: {})
    completed_steps << { step: step_name, result: result, at: Time.current.iso8601 }
    save!
  end
 
  def fail!(reason:, compensation_needed: [])
    update!(
      status: 'compensating',
      failure_reason: reason,
      compensation_steps: compensation_needed
    )
  end
end

どちらを選ぶか

Choreography を選ぶとき:
  - サービス数が少ない(3〜4サービス)
  - フローがシンプル
  - チームがイベント駆動に慣れている

Orchestration を選ぶとき:
  - サービス数が多い(5サービス以上)
  - 複雑な条件分岐がある
  - 可視性・デバッグのしやすさが重要
  - 補償フローが複雑

INFO

ShopNovaでは、最初はChoreographyで実装し、フローが複雑になるにつれてStep Functionsに移行した。これが現実的なアプローチ。最初から完璧な設計を目指さない。


まとめ

「分散トランザクションは『全部成功か全部失敗』ではなく、『失敗したら補償する』という考え方に変わる」とミサキが総括した。

サーガパターンの鉄則:
  1. 各ステップは冪等に設計する
  2. 各ステップは補償トランザクションを持つ
  3. 最終的な整合性を受け入れる
  4. saga_id でトレースできるようにする

次章では、サーガパターンで使ったイベント駆動アーキテクチャをさらに深掘りする。