mybook

Kata: リアルタイムチャット — WebSocket と非同期

課題の提示

EC Kata の翌週、ナオミがホワイトボードに「200ms以下」と書いた。

「これが今回のキーワード。ユーザーがメッセージを送ってから、相手に届くまでの時間。なぜこれが重要かわかる?」


Kata 3: リアルタイムチャットシステム

社内コミュニケーションツールのチャット機能を作ってほしい。

  • 同時接続ユーザー: 最大5,000人
  • メッセージ配信遅延: 200ms以下
  • チャンネル数: 最大1,000チャンネル
  • メッセージ履歴: 全て保存(永続化)
  • 既読機能がある
  • オンライン状態の表示
  • ファイル送信(画像・PDF)

タクミはすぐ気づいた。「HTTP リクエストじゃ無理ですよね。ポーリングは?」

ナオミは首を振った。「5,000人が1秒ごとにポーリングしたら?」

「5,000リクエスト/秒...それはまずい」

「そう。だから接続を"持続"させる仕組みが必要。WebSocket よ」


設計判断

HTTP vs WebSocket vs SSE

Loading diagram...

スケールの問題

WebSocket には HTTP にない課題がある。接続がサーバーに固定される問題だ。

Loading diagram...

解決策: パブサブ(Pub/Sub)バスでサーバー間を繋ぐ。

Loading diagram...

INFO

Rails の Action Cable は Redis を Pub/Sub バスとして使い、複数サーバー間でメッセージを同期する仕組みを内包している。


実装

Action Cable の設定

# config/cable.yml
development:
  adapter: async
 
production:
  adapter: redis
  url: <%= ENV.fetch("REDIS_URL") { "redis://localhost:6379/1" } %>
  channel_prefix: chat_app_production

チャンネルの実装

# app/channels/room_channel.rb
class RoomChannel < ApplicationCable::Channel
  def subscribed
    room = Room.find(params[:room_id])
    reject unless current_user.member_of?(room)
 
    stream_from "room:#{room.id}"
 
    # オンライン状態を更新
    PresenceService.mark_online(current_user, room)
    broadcast_presence(room, current_user, :online)
  end
 
  def unsubscribed
    room = Room.find_by(id: params[:room_id])
    return unless room
 
    PresenceService.mark_offline(current_user, room)
    broadcast_presence(room, current_user, :offline)
  end
 
  def receive(data)
    room = Room.find(params[:room_id])
    return unless current_user.member_of?(room)
 
    # バリデーション
    message_params = data.slice("content", "file_url")
    return if message_params["content"].blank? && message_params["file_url"].blank?
 
    # 非同期でメッセージ処理(ここでブロックしない)
    MessageBroadcastJob.perform_later(
      room_id: room.id,
      user_id: current_user.id,
      content: message_params["content"],
      file_url: message_params["file_url"]
    )
  end
 
  private
 
  def broadcast_presence(room, user, status)
    ActionCable.server.broadcast(
      "room:#{room.id}",
      { type: "presence", user_id: user.id, status: status }
    )
  end
end

メッセージのブロードキャスト (非同期Job)

# app/jobs/message_broadcast_job.rb
class MessageBroadcastJob < ApplicationJob
  queue_as :messages
 
  def perform(room_id:, user_id:, content:, file_url: nil)
    room = Room.find(room_id)
    user = User.find(user_id)
 
    # DBに保存
    message = room.messages.create!(
      user: user,
      content: content,
      file_url: file_url
    )
 
    # 全購読者にブロードキャスト
    ActionCable.server.broadcast(
      "room:#{room_id}",
      {
        type: "message",
        message: {
          id: message.id,
          content: message.content,
          file_url: message.file_url,
          created_at: message.created_at.iso8601,
          user: {
            id: user.id,
            name: user.name,
            avatar_url: user.avatar_url
          }
        }
      }
    )
  end
end

既読管理

既読は書き込み頻度が非常に高い(全メッセージ × 全ユーザー)。工夫が必要だ。

# app/models/message_read.rb
class MessageRead < ApplicationRecord
  belongs_to :user
  belongs_to :message
 
  # 既読はRedisで管理してバッチでDBに書く
  def self.mark_read_async(user_id, message_ids)
    # Redisに一時保存
    redis_key = "read_queue:#{user_id}"
    $redis.sadd(redis_key, message_ids)
    $redis.expire(redis_key, 1.hour)
 
    # バックグラウンドジョブで非同期にDB書き込み
    FlushReadStatusJob.perform_later(user_id)
  end
end
 
# app/jobs/flush_read_status_job.rb
class FlushReadStatusJob < ApplicationJob
  queue_as :low_priority
 
  def perform(user_id)
    redis_key = "read_queue:#{user_id}"
    message_ids = $redis.smembers(redis_key)
    return if message_ids.empty?
 
    # バルクインサートでDB効率化
    records = message_ids.map do |message_id|
      { user_id: user_id, message_id: message_id, created_at: Time.current }
    end
    MessageRead.upsert_all(records, unique_by: [:user_id, :message_id])
    $redis.del(redis_key)
  end
end

WARNING

既読の書き込みを全てリアルタイムにDBに書くと、5,000ユーザー × 1,000メッセージ/日 = 500万件/日のINSERTが発生する。Redis を一時バッファにしてバッチ処理するのが現実的。

オンライン状態管理

# app/services/presence_service.rb
class PresenceService
  ONLINE_TTL = 30.seconds
 
  def self.mark_online(user, room)
    key = "online:room:#{room.id}:user:#{user.id}"
    $redis.setex(key, ONLINE_TTL, 1)
 
    # ソート済みセット(オンラインユーザー一覧)
    $redis.zadd("online_users:room:#{room.id}", Time.current.to_i, user.id)
  end
 
  def self.mark_offline(user, room)
    $redis.del("online:room:#{room.id}:user:#{user.id}")
    $redis.zrem("online_users:room:#{room.id}", user.id)
  end
 
  def self.online_user_ids(room)
    cutoff = (Time.current - ONLINE_TTL).to_i
    $redis.zrangebyscore("online_users:room:#{room.id}", cutoff, "+inf").map(&:to_i)
  end
end

AWSインフラ構成

リアルタイム通信はインフラへの要求が高い。

Loading diagram...

ElastiCache Redis の設定

# Action Cable の Redis 接続設定
production:
  adapter: redis
  url: redis://chat-redis.xxxxx.clustercfg.use1.cache.amazonaws.com:6379
  channel_prefix: chat_production
 
# Redis Cluster の注意点:
# - 単一ノードより高可用性
# - channel_prefix でアプリ間の名前空間を分ける
# - connection_pool を設定して接続数を管理

INFO

ElastiCache Redis のクラスターモードは、最大数百GBのメモリと数十万の接続に対応できる。5,000ユーザー規模なら単一ノード(cache.r6g.large)で十分だが、将来のスケールに備えてクラスター構成にしておく。

ALB の WebSocket 設定

{
  "TargetGroupAttributes": [
    {
      "Key": "stickiness.enabled",
      "Value": "true"
    },
    {
      "Key": "stickiness.type",
      "Value": "lb_cookie"
    },
    {
      "Key": "stickiness.lb_cookie.duration_seconds",
      "Value": "86400"
    }
  ]
}

ALBのスティッキーセッション設定が必要だ。WebSocket 接続は TCP ベースで長時間維持されるため、同じユーザーの接続が常に同じサーバーに届く必要がある(Redis の Pub/Sub があれば必須ではないが、効率が上がる)。


パフォーマンスチューニング

Puma スレッド数の最適化

# config/puma.rb
# WebSocket は IO待ちが多い → スレッドを増やす
max_threads_count = ENV.fetch("RAILS_MAX_THREADS") { 10 }
min_threads_count = ENV.fetch("RAILS_MIN_THREADS") { max_threads_count }
threads min_threads_count, max_threads_count
 
workers ENV.fetch("WEB_CONCURRENCY") { 2 }

メッセージ履歴の最適化

# app/models/message.rb
class Message < ApplicationRecord
  belongs_to :room
  belongs_to :user
  has_many :reads, class_name: "MessageRead"
 
  # カーソルベースのページネーション(IDベース)
  scope :before, ->(id) { where("id < ?", id).order(id: :desc) }
 
  def self.history(room_id, before_id: nil, limit: 50)
    scope = where(room_id: room_id).includes(:user)
    scope = scope.before(before_id) if before_id
    scope.limit(limit).to_a.reverse
  end
end

INFO

チャットのページネーションは offset ではなくカーソル(ID)ベースが適切。新しいメッセージが到着するたびに offset がずれて「重複」や「欠け」が発生するのを防ぐ。


振り返り

タクミは実装を終えて気づいた。「WebSocket は難しくないが、スケールさせるのが難しい」

ナオミは言った。「正確に言うと、スケールではなく状態の共有が難しい。HTTP はステートレスだから、どのサーバーが処理しても同じ結果になる。WebSocket は接続に状態がある。その状態をどう共有するか」

「Redis の Pub/Sub ですね」

「そう。Action Cable はそれを隠蔽してくれる。でも隠蔽されているということは、内部で何が起きているか知らないと、デバッグできない」

INFO

Kata 3 の学び: リアルタイム通信は「状態の共有」が本質的な課題。WebSocket の採用は技術選択ではなく、その課題を解決する手段。Action Cable + Redis がその答えの一つ。

トレードオフの記録

決定メリットデメリット
Action CableRails統合、設定少Redis必須、スケール上限あり
既読のRedisバッファ書き込み負荷軽減瞬間的な不整合可能性
カーソルベースページネーション正確な履歴取得実装がoffsetより複雑

「次の Kata は予約システム。リアルタイムより"確実さ"が求められる世界よ」