mybook

Stage 10: 分散システム — 大規模設計の入り口

スケールの壁

入社から2年が経った。ヒロシのチームが作ったECサイトはヒットし、ユーザー数が10万人を超えた。

ある日、DBのCPUが100%に張り付き、サイトが数分間ダウンした。

「マイさん、どうすればよかったんでしょう。」

「単一のDBに全ての負荷をかけていたから。分散システムの考え方が必要なタイミングだった。」

「分散システムとは、複数のコンピュータが協調して一つのシステムとして動作すること。スケールの問題を解決するが、新しい問題も生む。」

CAP定理: 分散システムの基本制約

「まずCAP定理から。これは分散システムを語るとき避けられない。」

Loading diagram...

「現実のシステムでネットワーク分断(P)は避けられない。だからCPかAPかを選ぶことになる。」

選択特性向いているシステム
CP(一貫性優先)データは正確だが一時停止する銀行取引、在庫管理
AP(可用性優先)常に動くが一時的に古いデータSNSのいいね数、カート
# CP選択: 在庫管理(在庫の正確さが命)
class InventoryService
  def reserve!(product_id, quantity)
    # ペシミスティックロックで一貫性を保証
    product = Product.lock.find(product_id)
    raise InsufficientStock if product.stock < quantity
    product.decrement!(:stock, quantity)
  end
end
 
# AP選択: いいね数(厳密な正確さより応答速度優先)
class LikeCounter
  def increment(post_id)
    # Redisに書いて、定期的にDBに反映(最終的一貫性)
    $redis.incr("likes:#{post_id}")
    SyncLikesToDbJob.perform_later(post_id)
  end
 
  def count(post_id)
    $redis.get("likes:#{post_id}").to_i
  end
end

水平スケーリング: 台数で解決する

「単一サーバーの限界に達したとき、どうする?」

Loading diagram...
# ECS Auto Scaling設定
Resources:
  AppService:
    Type: AWS::ECS::Service
    Properties:
      ServiceName: myapp
      Cluster: !Ref Cluster
      TaskDefinition: !Ref TaskDefinition
      DesiredCount: 2  # 最低2台(可用性のため)
 
  AppScalingTarget:
    Type: AWS::ApplicationAutoScaling::ScalableTarget
    Properties:
      MaxCapacity: 20
      MinCapacity: 2
      ResourceId: !Sub "service/${Cluster}/${AppService.Name}"
      ScalableDimension: ecs:service:DesiredCount
      ServiceNamespace: ecs
 
  CPUScalingPolicy:
    Type: AWS::ApplicationAutoScaling::ScalingPolicy
    Properties:
      PolicyName: CPUScaling
      PolicyType: TargetTrackingScaling
      TargetTrackingScalingPolicyConfiguration:
        TargetValue: 70.0  # CPU 70%でスケールアウト
        PredefinedMetricSpecification:
          PredefinedMetricType: ECSServiceAverageCPUUtilization
# Railsをステートレスに保つ(水平スケーリングの前提)
# セッションはDBやRedisに保存
Rails.application.config.session_store :redis_store,
  servers: ENV.fetch("REDIS_URL"),
  expire_after: 90.minutes,
  key: "_app_session",
  secure: Rails.env.production?
 
# ファイルはS3に保存(ローカルファイルシステムに依存しない)
config.active_storage.service = :amazon

メッセージングと非同期処理

「全ての処理を同期的にやる必要はない。非同期でいい処理は非同期に。」

# 同期処理の問題: 全部完了するまでレスポンスを返せない
class OrdersController < ApplicationController
  def create
    order = Order.create!(order_params)
    OrderMailer.confirmation(order).deliver_now    # 2秒
    SlackNotifier.notify_new_order(order)          # 1秒
    InventoryService.update_stock(order)           # 3秒
    CrmService.sync_order(order)                   # 5秒
    render json: order  # 合計11秒後にレスポンス
  end
end
 
# 非同期処理: すぐにレスポンスして、残りをバックグラウンドで
class OrdersController < ApplicationController
  def create
    order = Order.create!(order_params)
 
    # Sidekiq/ActiveJobで非同期実行
    OrderConfirmationJob.perform_later(order.id)
    OrderSlackNotificationJob.perform_later(order.id)
    InventoryUpdateJob.perform_later(order.id)
    CrmSyncJob.perform_later(order.id)
 
    render json: order, status: :created  # 即座にレスポンス
  end
end
 
# app/jobs/order_confirmation_job.rb
class OrderConfirmationJob < ApplicationJob
  queue_as :default
  retry_on StandardError, wait: :polynomially_longer, attempts: 5
 
  def perform(order_id)
    order = Order.find(order_id)
    OrderMailer.confirmation(order).deliver_now
  end
end

SQSを使ったマイクロサービス間通信

# AWS SQSを使ったメッセージング
class OrderEventPublisher
  QUEUE_URL = ENV.fetch("ORDER_EVENTS_QUEUE_URL")
 
  def self.publish(event_type, payload)
    sqs = Aws::SQS::Client.new(region: "ap-northeast-1")
    sqs.send_message(
      queue_url: QUEUE_URL,
      message_body: {
        event_type: event_type,
        payload: payload,
        published_at: Time.current.iso8601,
        idempotency_key: SecureRandom.uuid
      }.to_json
    )
  end
end
 
# 注文作成時にSQSにイベントを発行
class Order < ApplicationRecord
  after_create_commit :publish_created_event
 
  private
 
  def publish_created_event
    OrderEventPublisher.publish("order.created", {
      order_id: id,
      user_id: user_id,
      total_amount: total_amount
    })
  end
end
 
# 在庫サービス側でイベントを消費
class InventoryEventConsumer
  def process(message)
    body = JSON.parse(message.body)
    case body["event_type"]
    when "order.created"
      InventoryReservationService.new(body["payload"]).call
    end
  end
end

耐障害性: 壊れることを前提に設計する

「分散システムは壊れる。ネットワークは切れる。サービスはダウンする。それを前提に設計するのが耐障害性設計。」

サーキットブレーカー

# 外部サービスが落ちているとき、無駄なリクエストを送らない
class CircuitBreaker
  FAILURE_THRESHOLD = 5
  RESET_TIMEOUT = 60  # 秒
 
  def initialize(name)
    @name = name
    @failures = 0
    @state = :closed  # :closed, :open, :half_open
    @last_failure_at = nil
  end
 
  def call
    case @state
    when :open
      if Time.current - @last_failure_at > RESET_TIMEOUT
        @state = :half_open
        attempt_call { yield }
      else
        raise CircuitOpenError, "サーキット #{@name} はオープン状態です"
      end
    when :closed, :half_open
      attempt_call { yield }
    end
  end
 
  private
 
  def attempt_call
    result = yield
    @failures = 0
    @state = :closed
    result
  rescue StandardError => e
    @failures += 1
    @last_failure_at = Time.current
    @state = :open if @failures >= FAILURE_THRESHOLD
    raise e
  end
end
 
# 使用例
class PaymentService
  def initialize
    @circuit_breaker = CircuitBreaker.new("stripe")
  end
 
  def charge(amount, token)
    @circuit_breaker.call do
      Stripe::Charge.create(amount: amount, source: token)
    end
  rescue CircuitOpenError
    { error: "決済サービスが一時的に利用できません", retry_after: 60 }
  end
end

リトライとべき等性

# べき等性: 同じリクエストを複数回送っても結果が同じ
class PaymentProcessor
  def process(idempotency_key:, amount:, token:)
    # 既に処理済みかチェック
    existing = Payment.find_by(idempotency_key: idempotency_key)
    return existing if existing
 
    # 初回のみ処理
    payment = nil
    ActiveRecord::Base.transaction do
      payment = Payment.create!(
        idempotency_key: idempotency_key,
        amount: amount,
        status: "pending"
      )
 
      charge = Stripe::Charge.create(amount: amount, source: token)
      payment.update!(stripe_charge_id: charge.id, status: "completed")
    end
 
    payment
  end
end

INFO

べき等性の重要性

ネットワークエラーでリトライが起きたとき、同じ処理が二重に実行されると困るケース(二重課金など)があります。idempotency_keyを使って「同じリクエストは一度だけ」を保証しましょう。

キャッシュ戦略

Loading diagram...
# Rails + ElastiCache (Redis) のキャッシュ戦略
 
# config/environments/production.rb
config.cache_store = :redis_cache_store, {
  url: ENV.fetch("REDIS_URL"),
  expires_in: 1.hour
}
 
# モデルレベルのキャッシュ
class Product < ApplicationRecord
  def self.featured
    Rails.cache.fetch("products/featured", expires_in: 15.minutes) do
      where(featured: true).includes(:category).to_a
    end
  end
 
  # キャッシュの無効化
  after_save :invalidate_featured_cache
 
  private
 
  def invalidate_featured_cache
    Rails.cache.delete("products/featured") if saved_change_to_featured?
  end
end
 
# HTTPキャッシュ(ETagとLast-Modified)
class ProductsController < ApplicationController
  def show
    @product = Product.find(params[:id])
 
    # ETagが一致すれば304 Not Modifiedを返す
    if stale?(etag: @product, last_modified: @product.updated_at)
      render json: @product
    end
  end
end

AWS での分散アーキテクチャ

# マイクロサービス的な構成
# 受注サービス → SQS → 在庫サービス
#                    → 通知サービス
#                    → 分析サービス
 
Services:
  OrderService:
    Type: ECS Service
    # 注文の受付・管理
 
  InventoryService:
    Type: ECS Service
    # 在庫の管理
 
  NotificationService:
    Type: Lambda
    # 通知の送信(イベント駆動)
 
  AnalyticsService:
    Type: Kinesis Data Firehose → S3 → Athena
    # イベントの分析
 
Message Bus:
  OrderEvents:
    Type: SNS Topic → SQS Queues
    # 各サービスが独立してスケール

Stage 10 のまとめ

「分散システムは銀の弾丸じゃない。複雑さが増す。」マイが言った。

「でも、スケールの要求がある場合は避けられない。大切なのは必要になってから導入すること。最初からマイクロサービスにする必要はない。」

「モノリス→モジュラーモノリス→マイクロサービスの順に、必要に応じて移行する。」

問題解法AWSサービス
単一サーバーの限界水平スケーリングECS Auto Scaling
同期処理のボトルネック非同期メッセージングSQS, EventBridge
外部サービス障害への連鎖サーキットブレーカー自前実装 or Hystrix
読み取り負荷キャッシュElastiCache Redis
データ整合性べき等性・最終的一貫性設計で対応

「次が最後のチャプター。エピローグだ。ヒロシがこれからどう成長し続けるか、その道標を話そう。」

ヒロシはDBのCPUが100%になった日の恐怖を思い出しながら、分散システムの学習をノートにまとめた。