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
endSQSを使ったマイクロサービス間通信
# 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
endINFO
べき等性の重要性
ネットワークエラーでリトライが起きたとき、同じ処理が二重に実行されると困るケース(二重課金など)があります。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
endAWS での分散アーキテクチャ
# マイクロサービス的な構成
# 受注サービス → 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%になった日の恐怖を思い出しながら、分散システムの学習をノートにまとめた。