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
スケールの問題
WebSocket には HTTP にない課題がある。接続がサーバーに固定される問題だ。
解決策: パブサブ(Pub/Sub)バスでサーバー間を繋ぐ。
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
endWARNING
既読の書き込みを全てリアルタイムに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
endAWSインフラ構成
リアルタイム通信はインフラへの要求が高い。
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
endINFO
チャットのページネーションは offset ではなくカーソル(ID)ベースが適切。新しいメッセージが到着するたびに offset がずれて「重複」や「欠け」が発生するのを防ぐ。
振り返り
タクミは実装を終えて気づいた。「WebSocket は難しくないが、スケールさせるのが難しい」
ナオミは言った。「正確に言うと、スケールではなく状態の共有が難しい。HTTP はステートレスだから、どのサーバーが処理しても同じ結果になる。WebSocket は接続に状態がある。その状態をどう共有するか」
「Redis の Pub/Sub ですね」
「そう。Action Cable はそれを隠蔽してくれる。でも隠蔽されているということは、内部で何が起きているか知らないと、デバッグできない」
INFO
Kata 3 の学び: リアルタイム通信は「状態の共有」が本質的な課題。WebSocket の採用は技術選択ではなく、その課題を解決する手段。Action Cable + Redis がその答えの一つ。
トレードオフの記録
| 決定 | メリット | デメリット |
|---|---|---|
| Action Cable | Rails統合、設定少 | Redis必須、スケール上限あり |
| 既読のRedisバッファ | 書き込み負荷軽減 | 瞬間的な不整合可能性 |
| カーソルベースページネーション | 正確な履歴取得 | 実装がoffsetより複雑 |
「次の Kata は予約システム。リアルタイムより"確実さ"が求められる世界よ」