サーガパターン — 分散トランザクション
「注文が確定したとき、在庫を減らして、決済を処理して、ポイントを付与する — これって全部成功か全部失敗にしたいですよね」
ミサキがうなずく前に、シニアエンジニアのアオイが言葉を続けた。「でも、それぞれ別サービスだから、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
endChoreography の問題点
メリット:
✓ シンプル、中央コントローラが不要
✓ 各サービスが疎結合
デメリット:
✗ 全体の流れが追いにくい(どのサービスが何を担当?)
✗ 循環依存が発生しやすい
✗ デバッグが困難(障害の原因を追うのが大変)
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' }
endRails から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 でトレースできるようにする
次章では、サーガパターンで使ったイベント駆動アーキテクチャをさらに深掘りする。