設計問題: チャットシステム — LINEのメッセージング設計
「リアルタイムの難しさ」
「今日のテーマはチャットシステム」とレイカが言った。「LINEのようなメッセージングアプリ。」
「ポーリングで解決できませんか?」ソウタが聞いた。
「できなくはない。でも1秒ごとにサーバーに問い合わせるポーリングを1億ユーザーがやったらどうなる?」
「……1億req/secですね。」
「そう。だからWebSocketが必要になる。まずHTTPとWebSocketを比較するところから始めよう。」
WebSocket vs HTTP Long Polling
HTTP ポーリング(非推奨)
# クライアント側の擬似コード(ポーリング)
def polling_loop
loop do
response = HTTP.get("/messages?since=#{last_timestamp}")
messages = response.json
if messages.any?
display(messages)
self.last_timestamp = messages.last[:timestamp]
end
sleep(1) # 1秒ごとに確認
end
end
# 問題点:
# 1. 1秒以内のリアルタイム性が保証できない
# 2. 1億ユーザー × 1req/sec = 1億req/sec(サーバーが壊れる)
# 3. 「新着なし」のレスポンスが99%を占める(無駄)HTTP Long Polling(ポーリングの改善版)
# Long Polling: サーバーが新着メッセージを「待って」から返す
class MessagesController < ApplicationController
def poll
# 最大30秒待機、メッセージが来たら即座に返す
message = Message.wait_for_new(
room_id: params[:room_id],
since: params[:since],
timeout: 30
)
if message
render json: message
else
render json: { status: 'timeout' }, status: :no_content
end
end
end
# 問題点:
# 1. サーバーコネクションが長時間占有される
# 2. 1億ユーザーが30秒待つ = 1億本のコネクション
# 3. HTTPオーバーヘッドが毎回発生(ヘッダー数KB)WebSocket(推奨)
WebSocketのメリット:
1. 一度接続を確立したら維持(再ネゴシエーション不要)
2. 双方向リアルタイム通信
3. サーバーからクライアントへのPush配信が可能
4. HTTPオーバーヘッドが初回のみ(その後は数バイトのフレーム)
接続確立フロー:
Client: GET /cable HTTP/1.1
Upgrade: websocket
Connection: Upgrade
Server: HTTP/1.1 101 Switching Protocols
Upgrade: websocket
→ 以降はバイナリフレームで双方向通信
# Rails ActionCableでのWebSocket実装
# app/channels/chat_channel.rb
class ChatChannel < ApplicationCable::Channel
def subscribed
room = ChatRoom.find(params[:room_id])
reject unless current_user.member_of?(room)
stream_for room
current_user.update_online_status(online: true)
broadcast_presence(room, :online)
end
def unsubscribed
current_user.update_online_status(online: false)
# 所属する全ルームにオフライン通知
current_user.chat_rooms.each do |room|
broadcast_presence(room, :offline)
end
end
def send_message(data)
room = ChatRoom.find(data['room_id'])
return reject unless current_user.member_of?(room)
message = Message.create!(
sender: current_user,
chat_room: room,
content: data['content'],
id: SnowflakeGenerator.next_id
)
ChatChannel.broadcast_to(room, {
type: 'new_message',
message: MessageSerializer.new(message).as_json
})
PushNotificationJob.perform_later(message.id)
end
private
def broadcast_presence(room, status)
ChatChannel.broadcast_to(room, {
type: 'presence',
user_id: current_user.id,
status: status
})
end
endStep 1: 要件の確認
ソウタ:「確認させてください。
【機能要件】
- 1対1のプライベートチャット
- グループチャット(最大500人)
- テキストメッセージ(画像は後で拡張)
- メッセージの既読確認(既読/未読)
- オンライン/オフラインステータス
- プッシュ通知(オフライン時)
- @メンション機能
【非機能要件】
- DAU: 5,000万(50M)
- 同時接続数: 1,000万(10M)
- メッセージのレイテンシ: P99で100ms以内
- メッセージ配信保証: At-least-once
- 可用性: 99.99%」
Step 2: 規模の概算
【QPS計算】
同時接続数: 10M
1ユーザーあたりの1日のメッセージ数: 20件
メッセージQPS: 50M × 20 / 86,400 ≒ 11,600 messages/sec
ピーク: 35,000 messages/sec
【ストレージ】
メッセージ1件: sender_id(8B) + receiver_id(8B) + content(1KB) + ts(8B) ≒ 1KB
1日のデータ量: 50M × 20 × 1KB = 1TB/day
5年間: 1TB × 365 × 5 ≒ 1.8PB
【WebSocket接続数】
同時接続: 10M
Chat Server 1台での管理: 約10万〜100万(チューニング次第)
必要なChat Server数: 10M / 100K = 100台
Step 3: 高レベル設計
Step 4: ステートレス設計 — どのサーバーに繋いでも同じ体験
「100台のChat Serverがある場合、ClientAとClientBは別々のサーバーに繋がっている可能性が高い。どう解決する?」とレイカが聞いた。
# app/services/connection_registry.rb
# どのChat Serverにどのユーザーが接続しているかをRedisで管理
class ConnectionRegistry
def self.register(user_id, server_id)
$redis.setex("user_server:#{user_id}", 300, server_id) # 5分TTL
end
def self.find_server(user_id)
$redis.get("user_server:#{user_id}")
end
def self.unregister(user_id)
$redis.del("user_server:#{user_id}")
end
end
# メッセージルーティング(Redis Pub/Sub)
class MessageRouter
def self.deliver(message)
receiver_id = message.receiver_id
server_id = ConnectionRegistry.find_server(receiver_id)
if server_id
# オンライン: 該当Chat Serverに転送
$redis.publish("server:#{server_id}", {
type: 'message',
payload: MessageSerializer.new(message).as_json
}.to_json)
else
# オフライン: DBに保存してPush通知
store_offline_message(message)
PushNotificationJob.perform_later(message.id)
end
end
endStep 5: メッセージIDの設計(Snowflake)
「チャットシステムで意外に難しいのがメッセージIDだ」とレイカが言った。
要件:
- グローバルにユニークであること
- 時系列でソート可能であること(メッセージの順序保証)
- 生成が高速であること(35,000 msg/sec に対応)
方式比較:
UUID v4: ランダムなため、時系列ソートができない ❌
DB AUTO_INCREMENT: 単一DBでは良いが、分散DBでは衝突 ❌
Snowflake形式: 時刻ベース + ノードID + シーケンス番号 ✅
# app/services/snowflake_generator.rb
class SnowflakeGenerator
EPOCH = 1_700_000_000_000 # 2023-11-15基準(ms)
NODE_BITS = 10
SEQ_BITS = 12
MAX_SEQ = (1 << SEQ_BITS) - 1 # 4095
@@sequence = 0
@@last_ts = -1
@@mutex = Mutex.new
def self.next_id
@@mutex.synchronize do
ts = current_ms
if ts == @@last_ts
@@sequence = (@@sequence + 1) & MAX_SEQ
ts = wait_next_ms if @@sequence == 0
else
@@sequence = 0
end
@@last_ts = ts
((ts - EPOCH) << (NODE_BITS + SEQ_BITS)) |
(node_id << SEQ_BITS) |
@@sequence
end
end
private
def self.current_ms = (Time.now.to_f * 1000).to_i
def self.node_id = ENV.fetch('NODE_ID').to_i
def self.wait_next_ms
ts = current_ms
ts <= @@last_ts ? wait_next_ms : ts
end
endStep 6: Golangでのチャットサーバー実装
「Golangのgoroutineとchannelはチャットサーバーと非常に相性が良い」とレイカが言った。
// internal/chat/server.go
package chat
import (
"context"
"encoding/json"
"sync"
"github.com/gorilla/websocket"
)
type Message struct {
ID int64 `json:"id"`
SenderID int64 `json:"sender_id"`
ReceiverID int64 `json:"receiver_id"`
Content string `json:"content"`
Timestamp int64 `json:"timestamp"`
}
type Client struct {
UserID int64
conn *websocket.Conn
send chan []byte // 送信バッファ
}
type Hub struct {
mu sync.RWMutex
clients map[int64]*Client // userID → Client
register chan *Client
unregister chan *Client
broadcast chan Message
}
func NewHub() *Hub {
return &Hub{
clients: make(map[int64]*Client),
register: make(chan *Client),
unregister: make(chan *Client),
broadcast: make(chan Message, 1000),
}
}
func (h *Hub) Run(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case client := <-h.register:
h.mu.Lock()
h.clients[client.UserID] = client
h.mu.Unlock()
case client := <-h.unregister:
h.mu.Lock()
if _, ok := h.clients[client.UserID]; ok {
delete(h.clients, client.UserID)
close(client.send)
}
h.mu.Unlock()
case msg := <-h.broadcast:
h.mu.RLock()
client, ok := h.clients[msg.ReceiverID]
h.mu.RUnlock()
if ok {
data, _ := json.Marshal(msg)
select {
case client.send <- data:
default:
// バッファが詰まっている場合は接続を切る
h.mu.Lock()
delete(h.clients, client.UserID)
close(client.send)
h.mu.Unlock()
}
} else {
// オフライン: Push通知を送る
go SendPushNotification(msg)
}
}
}
}
// writePumpは各クライアントごとのgoroutine
func (c *Client) writePump() {
defer c.conn.Close()
for {
msg, ok := <-c.send
if !ok {
c.conn.WriteMessage(websocket.CloseMessage, []byte{})
return
}
c.conn.WriteMessage(websocket.TextMessage, msg)
}
}Step 7: グループチャットの設計
-- グループチャットのデータモデル
CREATE TABLE chat_rooms (
id BIGINT PRIMARY KEY,
name VARCHAR(100),
room_type VARCHAR(20) NOT NULL, -- 'direct' or 'group'
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE TABLE room_members (
room_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
joined_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
last_read_message_id BIGINT, -- 既読位置
PRIMARY KEY (room_id, user_id)
);
CREATE INDEX idx_room_members_user ON room_members(user_id);
CREATE TABLE messages (
id BIGINT PRIMARY KEY, -- Snowflake ID
room_id BIGINT NOT NULL,
sender_id BIGINT NOT NULL,
content TEXT NOT NULL,
created_at TIMESTAMPTZ NOT NULL
);# app/services/group_message_broadcaster.rb
class GroupMessageBroadcaster
def self.broadcast(message)
room = ChatRoom.find(message.room_id)
members = RoomMember.where(room_id: room.id).pluck(:user_id)
online_ids = members.select { |uid| ConnectionRegistry.find_server(uid) }
offline_ids = members - online_ids
# オンラインメンバーへWebSocket配信
online_ids.each do |uid|
server_id = ConnectionRegistry.find_server(uid)
$redis.publish("server:#{server_id}", {
type: 'group_message',
payload: MessageSerializer.new(message).as_json
}.to_json)
end
# オフラインメンバーへPush通知
# バッジカウントをインクリメント
offline_ids.each do |uid|
$redis.incr("unread:#{uid}:#{room.id}")
PushNotificationJob.perform_later(message.id, uid)
end
end
endStep 8: メンション機能(@ユーザー名)の実装
# app/services/mention_parser.rb
class MentionParser
MENTION_PATTERN = /@(\w+)/
def self.parse_and_notify(message)
mentions = message.content.scan(MENTION_PATTERN).flatten
return if mentions.empty?
users = User.where(username: mentions)
users.each do |user|
# メンション通知を作成
Notification.create!(
user_id: user.id,
type: 'mention',
message_id: message.id,
room_id: message.room_id
)
# リアルタイム通知(オンラインの場合)
server_id = ConnectionRegistry.find_server(user.id)
if server_id
$redis.publish("server:#{server_id}", {
type: 'mention',
payload: {
message_id: message.id,
room_id: message.room_id,
from: message.sender.username
}
}.to_json)
end
end
end
end
# app/models/message.rb
class Message < ApplicationRecord
after_create :parse_mentions
private
def parse_mentions
MentionParser.parse_and_notify(self)
end
endStep 9: ファイル/画像送信の設計(S3 Presigned URL)
# app/controllers/api/v1/attachments_controller.rb
class Api::V1::AttachmentsController < ApplicationController
def presign
key = "attachments/#{current_user.id}/#{SecureRandom.uuid}/#{params[:filename]}"
expiration = 10.minutes
presigned_url = S3_CLIENT.presigned_url(
:put_object,
bucket: ENV['S3_BUCKET'],
key: key,
expires_in: expiration.to_i,
content_type: params[:content_type]
)
render json: {
presigned_url: presigned_url,
key: key,
expires_at: expiration.from_now
}
end
def confirm
# アップロード完了後、メッセージとして送信
message = Message.create!(
sender: current_user,
room_id: params[:room_id],
type: 'attachment',
content: params[:filename],
s3_key: params[:key],
file_size: params[:file_size]
)
GroupMessageBroadcaster.broadcast(message)
render json: MessageSerializer.new(message)
end
endStep 10: オフライン時のメッセージ保存と再同期
# app/jobs/offline_sync_job.rb
class OfflineSyncJob < ApplicationJob
# ユーザーがオンラインに戻った時に呼ばれる
def perform(user_id)
# 未読メッセージを取得(最後にオンラインだった時刻以降)
last_seen = UserStatus.find_by(user_id: user_id)&.last_seen_at || 7.days.ago
undelivered = Message.joins(:room_members)
.where(room_members: { user_id: user_id })
.where('messages.created_at > ?', last_seen)
.where('messages.id > room_members.last_read_message_id')
.order(:created_at)
if undelivered.any?
# WebSocketでまとめて送信
server_id = ConnectionRegistry.find_server(user_id)
if server_id
$redis.publish("server:#{server_id}", {
type: 'sync',
messages: MessageSerializer.new(undelivered).as_json
}.to_json)
end
end
# 未読バッジをリセット
$redis.del(*undelivered.map { |m| "unread:#{user_id}:#{m.room_id}" }.uniq)
end
endStep 11: メッセージの配信保証(At-least-once)
配信保証のレベル:
At-most-once: 失ってもOK(ログなど)
At-least-once: 重複あり可(今回採用)→ 受信側で重複除去
Exactly-once: 重複も欠損もない(最も複雑・コスト高)
At-least-onceの実装:
1. クライアントがメッセージを送信(ACK待ち)
2. サーバーが受信したらACKを返す
3. クライアントがACKを受け取れなかった場合、再送信
4. 受信側はSnowflake IDで重複チェック(冪等性)
# app/models/message.rb
class Message < ApplicationRecord
# べき等な作成(同じIDが来ても無視)
def self.create_idempotent!(params)
insert(params, unique_by: :id, returning: false)
rescue ActiveRecord::RecordNotUnique
find(params[:id]) # 重複の場合は既存を返す
end
end// クライアント側の再送ロジック(Golang)
func (c *ChatClient) SendWithRetry(ctx context.Context, msg Message) error {
for attempt := 0; attempt < 3; attempt++ {
if err := c.send(ctx, msg); err == nil {
return nil
}
waitMs := time.Duration(100*(1<<attempt)) * time.Millisecond
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(waitMs):
}
}
return fmt.Errorf("failed to send after 3 attempts")
}Step 12: 既読機能の設計
# db/schema.rb(既読管理)
# メッセージごとではなく「最後に読んだメッセージID」で管理する(書き込み削減)
create_table :room_members do |t|
t.bigint :room_id, null: false
t.bigint :user_id, null: false
t.bigint :last_read_message_id # これ以前は既読
t.datetime :created_at, null: false
end
# app/channels/chat_channel.rb(既読マーク)
def mark_as_read(data)
room_id = data['room_id']
message_id = data['message_id']
# 最後に読んだメッセージIDを更新
RoomMember.where(room_id: room_id, user_id: current_user.id)
.update_all(last_read_message_id: message_id)
# 送信者に既読を通知
ChatChannel.broadcast_to(ChatRoom.find(room_id), {
type: 'read_receipt',
message_id: message_id,
read_by: current_user.id,
read_at: Time.current.iso8601
})
# 未読バッジをクリア
$redis.del("unread:#{current_user.id}:#{room_id}")
endStep 13: データベースの選択(Cassandra)
PostgreSQLの問題:
- 35,000 msg/sec の大量書き込みにシングルマスターでは限界
- チャット履歴は時系列データ → テーブルが巨大になると検索が遅い
- 水平スケールが難しい
Cassandraが適している理由:
- 書き込み最適化(LSMツリー構造)
- 時系列データに相性の良いパーティション設計
- 水平スケール(ノード追加で自動分散)
- マルチリージョン対応
パーティション設計:
PARTITION KEY: room_id + year_month
CLUSTERING KEY: message_id DESC
→ "room_id=123の2024年11月のメッセージ一覧"が高速に取れる
-- CassandraのCQL定義
CREATE TABLE messages (
room_id BIGINT,
year_month TEXT,
message_id BIGINT,
sender_id BIGINT,
content TEXT,
msg_type TEXT, -- 'text', 'attachment', 'system'
created_at TIMESTAMP,
PRIMARY KEY ((room_id, year_month), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);
-- 効率的なクエリ例
-- room_id=123の2024年11月のメッセージを取得
SELECT * FROM messages
WHERE room_id = 123 AND year_month = '2024-11'
LIMIT 20;Step 14: AWSアーキテクチャの詳細
NLB(Network Load Balancer):
理由: WebSocketは長時間コネクションのためNLBが適切
ALBはHTTP/2までの対応でWebSocketも使えるが
NLBの方がTCPレイヤーで安定している
スティッキーセッション: 有効(同一クライアントを同一サーバーへ)
Chat Server(EC2):
インスタンス: c6i.4xlarge(16vCPU, 32GB)
台数: 100台(Auto Scaling)
1台あたりの接続数: 最大10万接続
理由: ECS Fargateではなく EC2(長時間接続の管理が容易)
Amazon MSK(Kafka):
ブローカー数: 3(マルチAZ)
Partition数: 100(並列処理)
Replication: 3
メッセージ保持期間: 7日間
Amazon Keyspaces(Cassandra互換):
代替: AWSマネージドのため運用負荷が低い
注意: Cassandraの一部機能は未サポート
代替案: EC2上のCassandraクラスター(より柔軟)
ElastiCache Redis:
用途:
- ConnectionRegistry(ユーザー→サーバーマッピング)
- オンラインステータス管理
- 未読バッジカウント
ノード: cache.r6g.xlarge × 3INFO
Amazon Keyspacesは、AWSマネージドのCassandra互換サービス。Cassandraのパーティション設計をそのまま使いながら、インフラ管理の手間を省ける。面接でCassandraの設計を語りつつ「AWSではKeyspacesを利用します」と言うと実用性が伝わる。
Step 15: 監視設計
# app/middleware/chat_metrics.rb
class ChatMetrics
WS_CONNECTIONS = Prometheus::Client::Gauge.new(
:websocket_connections_total,
docstring: 'Active WebSocket connections'
)
MSG_LATENCY = Prometheus::Client::Histogram.new(
:message_delivery_latency_ms,
docstring: 'Message delivery latency in milliseconds',
buckets: [5, 10, 25, 50, 100, 250, 500, 1000]
)
MSG_THROUGHPUT = Prometheus::Client::Counter.new(
:messages_sent_total, docstring: 'Total messages sent'
)
end# CloudWatch アラーム設定
アラーム一覧:
WebSocketConnectionDrop:
メトリクス: websocket_connections_total
閾値: 30秒で10%以上の急減
アクション: PagerDuty通知(サーバー障害の可能性)
MessageLatencyHigh:
メトリクス: p99(message_delivery_latency_ms)
閾値: 100ms超
アクション: PagerDuty通知
KafkaConsumerLag:
メトリクス: kafka.consumer_lag(MessageProcessorグループ)
閾値: 50,000メッセージ超
アクション: Message Processorのスケールアウト
KeyspacesThrottled:
メトリクス: Keyspaces SystemErrors
閾値: 100/min超
アクション: Capacity Units増加E2E暗号化の概要
「E2E暗号化(End-to-End Encryption)について聞かれたらどう答える?」とレイカが聞いた。
E2E暗号化の概念:
- メッセージはクライアント側で暗号化される
- サーバーは暗号化されたデータを転送するだけ
- サーバーはメッセージの内容を復号できない
- LINEの「Letter Sealing」、WhatsApp、Signal が採用
実装の核心(Signal Protocol):
1. 各ユーザーが公開鍵/秘密鍵のペアを生成
2. 公開鍵のみサーバーに登録
3. 送信者はAの公開鍵でメッセージを暗号化
4. 受信者はBの秘密鍵で復号
面接でのポジション:
「今回のスコープ外ですが、本番システムでは
Signal Protocolのような既存のE2E暗号化プロトコルを
採用することを推奨します。
キー管理の複雑さとUX(デバイス間の同期)がトレードオフです。」
面接官との深掘り会話
面接官: 「WebSocketサーバーが落ちたらどうなりますか?」
ソウタ: 「クライアントはWebSocket切断を検知して、
指数バックオフ(1秒→2秒→4秒→...)で再接続を試みます。
サーバーはステートレス設計にしてあるので、
どのChat Serverに再接続しても同じ体験が提供できます。
再接続後はOfflineSyncJobが未配信メッセージを送信します。」
---
面接官: 「グループチャット500人に同時配信するには?」
ソウタ: 「500人の接続状態をConnectionRegistryで確認し、
オンラインのメンバーが繋がっているサーバーに
Redis Pub/Subでメッセージを転送します。
オフラインのメンバーにはPush通知を送り、
未読バッジカウントをインクリメントします。
500人の場合、最悪500台のサーバーに転送が必要になりますが、
Redis Pub/Subはチャネルベースなので実装はシンプルです。」
---
面接官: 「メッセージの順序保証はどうしますか?」
ソウタ: 「Snowflake IDは時刻ベースのため、
同一サーバーで生成される限り単調増加が保証されます。
ただし分散環境では、異なるサーバーで1ms以内に
生成されたIDの順序が逆転する可能性があります。
クライアント側でIDでソートし直すか、
ルームごとにシーケンス番号を管理する方法があります。
LINEは後者のアプローチを取っています。」
WARNING
「エンドツーエンド暗号化(E2E暗号化)」は面接で触れると良いポイント。「サーバーはメッセージの内容を復号できない」という設計はLINEやWhatsAppが採用している。ただし実装は複雑なため、「スコープ外だが重要なセキュリティ要件です」と言える程度で良い。
ソウタの振り返り
「WebSocketとHTTPの比較から始めて、ステートレス設計、グループ配信、メッセージID、Cassandraのパーティション設計……全部つながっていますね。」
「チャットシステムは分散システムの難問が凝縮している」とレイカが言った。「リアルタイム配信、配信保証、順序保証、スケールアウト。どれか一つでも間違えると、ユーザーはメッセージが届かなかったと感じる。」
「一番難しかったのは、ステートレス設計です。どのサーバーに繋いでも同じ体験、というのはRedisを使えばシンプルに解決できるんですが、それに気づくまでが……」
「そこが設計の醍醐味だ。AnswerよりもThinking Processが評価される。どう考えてその結論に至ったかを、面接官は見ている。」
INFO
チャットシステムは面接で「深掘り耐性」が最も試される問題。WebSocket、Snowflake ID、Cassandraのパーティション設計、At-least-once配信、ステートレス設計の5つを体系立てて説明できると、高い評価を得られる。
次は動画配信プラットフォームの設計。YouTubeはどうやって数億本の動画を配信しているのか。