mybook

設計問題: SNSのタイムライン — Twitterのフィード設計

「プロローグの借り」

3ヶ月前のことをソウタはまだ覚えていた。第一志望の会社の最終面接。コーディングテストは満点だった。でも設計面接で「Twitterのタイムラインを設計してください」と聞かれた瞬間、頭が真っ白になった。

「ユーザーのツイートをDBに保存して、タイムライン取得時にフォロー中のユーザーのツイートをJOINで取ってくれば……」と答え始めた瞬間、面接官の表情が変わった。

「1億DAUでも同じですか?」

沈黙が続いた。結果は不採用だった。

「今日は、あなたが最初の面接で詰まった問題をやる」とレイカが言った。

「Twitterのタイムライン……ですね。」ソウタの声には特別な重みがあった。

「正確には、ホームタイムライン。自分がフォローしているユーザーの投稿を時系列で表示する機能。一見シンプルだけど、1億DAUになったとき、どうやってスケールさせるかが問われる。」

「今度こそ、完全に答えてみせます。」


Step 1: 要件の確認

ソウタ:「確認させてください。

【機能要件】
- ツイートの投稿(テキスト、最大280文字)
- ホームタイムラインの表示(フォロー中ユーザーのツイート)
- ユーザープロフィールタイムラインの表示
- フォロー/アンフォロー機能
- いいね・リツイート・画像は今回対象外

【非機能要件】
- DAU: 1億(100M)
- フォロワー数: 平均200人(最大: 有名人は数千万人)
- 書き込みQPS: 5,000/sec(ピーク: 15,000/sec)
- タイムライン取得: P99で100ms以内
- 可用性: 99.99%(月間ダウン4.3分以内)」

「要件確認だけで面接官はかなり評価している」とレイカが言った。「特に『フォロワーの最大数』を聞けるかどうか。これが設計全体を左右する。」

INFO

フォロワーの最大数を要件で確認することが最重要。一般ユーザーは200人だが、Elonのような有名人は1億人のフォロワーを持つ。この差が設計の核心「有名人問題」につながる。


Step 2: 規模の概算

【ツイート投稿QPS】
DAU: 100M
投稿するユーザー: 100M × 10% = 10M人/day
1人が平均5ツイート: 10M × 5 = 50M tweets/day
QPS: 50M / 86,400 ≒ 580/sec(ピーク3倍: 1,740/sec)

【タイムライン閲覧QPS】
1人が1日に20回タイムラインを開く
1回で20件のツイートを取得
QPS: 100M × 20 / 86,400 ≒ 23,000 req/sec

【ストレージ(ツイート本文)】
1件: tweet_id(8B) + user_id(8B) + text(280B) + created_at(8B) = 約300B
1日: 50M × 300B = 15 GB/day
10年: 15GB × 3,650 ≒ 55 TB

【比較】
書き込みQPS: 580/sec(ピーク1,740/sec)
読み込みQPS: 23,000/sec(ピーク70,000/sec)
→ 読み込みが40倍多い。読み込み最適化が設計の中心になる

Step 3: 高レベル設計とファンアウト

「タイムライン設計の核心はファンアウトの方式選択だ」とレイカが言った。

ファンアウトとは?

Loading diagram...

ファンアウトとは、ユーザーAがツイートしたとき、そのツイートをAのフォロワー全員に届ける処理。この「届け方」に3つの方式がある。

方式1: Push(書き込み時ファンアウト)

# ツイート投稿時 → 全フォロワーのタイムラインキャッシュに書き込む
class PushFanoutService
  def self.fanout(tweet)
    follower_ids = Follow.where(followee_id: tweet.user_id).pluck(:follower_id)
 
    follower_ids.each_slice(1000) do |batch|
      batch.each do |follower_id|
        $redis.lpush("timeline:#{follower_id}", tweet.id)
        $redis.ltrim("timeline:#{follower_id}", 0, 799)  # 最新800件のみ
      end
    end
  end
end
 
# メリット: 読み込みが高速(事前計算済み)
# デメリット: フォロワー100万人のユーザーが1ツイートで100万回書き込み

方式2: Pull(読み込み時ファンアウト)

# タイムライン取得時 → フォロー中ユーザーのツイートをDBから取得
class PullTimelineService
  def self.fetch(user_id, limit: 20)
    followee_ids = Follow.where(follower_id: user_id).pluck(:followee_id)
 
    # フォロー200人 × SELECT = 非効率なN+1に近い処理
    Tweet.where(user_id: followee_ids)
         .order(created_at: :desc)
         .limit(limit)
  end
end
 
# メリット: 書き込み時の負荷がない
# デメリット: フォロー200人のツイートを毎回結合するので遅い

方式3: ハイブリッド(Twitterの実際の方式)

INFO

Twitterは実際にハイブリッド方式を採用している。一般ユーザーはPush(事前にキャッシュへ書き込む)、フォロワーが多い「セレブリティ(有名人)」はPull(読み込み時にマージ)。この切り替えが「有名人問題(Celebrity Problem)」への解答だ。

# app/services/fanout_service.rb
class FanoutService
  CELEBRITY_THRESHOLD = 10_000  # フォロワー数の閾値
 
  def self.fanout(tweet)
    user = User.find(tweet.user_id)
 
    if user.followers_count < CELEBRITY_THRESHOLD
      # 一般ユーザー: Pushモード
      PushFanoutJob.perform_later(tweet.id)
    else
      # 有名人: Pullモード用にマーキングのみ
      $redis.sadd("celebrity_tweets", tweet.id)
      $redis.expire("celebrity_tweets", 7.days.to_i)
    end
  end
end

Step 4: Golangでの並列ファンアウト実装

「大規模なファンアウトはGolangのgoroutineで並列化するのが定石だ」とレイカが言った。

// internal/fanout/worker.go
package fanout
 
import (
	"context"
	"sync"
)
 
type FanoutWorker struct {
	redis     RedisClient
	batchSize int
}
 
func (w *FanoutWorker) Fanout(ctx context.Context, tweetID int64, followerIDs []int64) error {
	batches := chunk(followerIDs, w.batchSize) // 1000件ずつバッチ
 
	var wg sync.WaitGroup
	errCh := make(chan error, len(batches))
 
	for _, batch := range batches {
		wg.Add(1)
		go func(ids []int64) {
			defer wg.Done()
			if err := w.pushToTimelines(ctx, tweetID, ids); err != nil {
				errCh <- err
			}
		}(batch)
	}
 
	wg.Wait()
	close(errCh)
 
	// エラー収集
	for err := range errCh {
		if err != nil {
			return err
		}
	}
	return nil
}
 
func (w *FanoutWorker) pushToTimelines(ctx context.Context, tweetID int64, followerIDs []int64) error {
	pipe := w.redis.Pipeline()
	for _, followerID := range followerIDs {
		key := fmt.Sprintf("timeline:%d", followerID)
		pipe.LPush(ctx, key, tweetID)
		pipe.LTrim(ctx, key, 0, 799)
	}
	_, err := pipe.Exec(ctx)
	return err
}
 
func chunk(ids []int64, size int) [][]int64 {
	var chunks [][]int64
	for size < len(ids) {
		ids, chunks = ids[size:], append(chunks, ids[0:size:size])
	}
	return append(chunks, ids)
}

Step 5: Redisのデータ構造詳細

「RedisのListとSorted Setは何が違うか説明して」とレイカが聞いた。

【LPUSH + LTRIM(List)】
長所:
  - 追加がO(1)と高速
  - タイムラインのような「最新順」に最適
  - 実装がシンプル

短所:
  - 任意のスコアでソートできない
  - ランダムアクセスがO(N)

使いどころ:
  - タイムライン(時系列順は投稿順で決まる)

【ZADD(Sorted Set)】
長所:
  - スコア(タイムスタンプ)で自動ソート
  - 範囲取得がO(log N + M)
  - カーソルベースページネーションに最適

短所:
  - メモリ使用量がListより多い
  - 追加がO(log N)

使いどころ:
  - スコアリングが必要なランキング
  - カーソルが飛び回るページネーション
# Sorted Setを使ったタイムライン管理
class TimelineWithZSet
  def self.add_tweet(follower_id, tweet_id, score)
    key = "tl:zset:#{follower_id}"
    $redis.zadd(key, score, tweet_id)
    $redis.zremrangebyrank(key, 0, -801)  # 古い800件を超えたら削除
  end
 
  # カーソルベースページネーション
  def self.fetch(follower_id, cursor: nil, limit: 20)
    key   = "tl:zset:#{follower_id}"
    max   = cursor ? "(#{cursor}" : '+inf'  # カーソル未満のスコアを取得
    ids   = $redis.zrevrangebyscore(key, max, '-inf', limit: [0, limit])
    next_cursor = ids.last ? $redis.zscore(key, ids.last) : nil
    { tweet_ids: ids, next_cursor: next_cursor }
  end
end

カーソルベース vs オフセットベースのページネーション

# オフセットベース(シンプルだが問題がある)
def fetch_offset(page:, per: 20)
  # 問題: page=2 取得中に新しいツイートが追加されると
  #       1ページ目の最後のツイートが2ページ目に流れて重複する
  tweet_ids = $redis.lrange("timeline:#{user_id}", page * per, (page + 1) * per - 1)
end
 
# カーソルベース(推奨)
def fetch_cursor(cursor: nil, limit: 20)
  # カーソル = 前回取得した最後のtweetのスコア(タイムスタンプ)
  # 新しいツイートが追加されても、カーソル位置が変わらない
  result = TimelineWithZSet.fetch(user_id, cursor: cursor, limit: limit)
  {
    tweets:      Tweet.where(id: result[:tweet_ids]).index_by(&:id),
    next_cursor: result[:next_cursor]
  }
end

Step 6: 有名人問題の詳細解説

Loading diagram...
# app/services/timeline_service.rb
class TimelineService
  CELEBRITY_THRESHOLD = 10_000
 
  def self.get_home_timeline(user_id, cursor: nil, limit: 20)
    # 1. Pushキャッシュからの取得(一般ユーザー分)
    push_tweets = fetch_push_timeline(user_id, cursor: cursor, limit: limit * 2)
 
    # 2. フォロー中のセレブリティを特定
    celebrity_ids = Follow.where(follower_id: user_id)
                          .joins(:followee)
                          .where('users.followers_count >= ?', CELEBRITY_THRESHOLD)
                          .pluck(:followee_id)
 
    # 3. セレブリティのツイートをPull
    celebrity_tweets = []
    if celebrity_ids.any?
      since = cursor ? Time.at(cursor.to_f / 1000) : 48.hours.ago
      celebrity_tweets = Tweet.where(user_id: celebrity_ids)
                               .where('created_at > ?', since)
                               .order(created_at: :desc)
                               .limit(limit)
    end
 
    # 4. マージして時系列ソート
    all_tweets = (push_tweets + celebrity_tweets)
                   .uniq(&:id)
                   .sort_by(&:created_at)
                   .reverse
                   .first(limit)
 
    all_tweets
  end
 
  private
 
  def self.fetch_push_timeline(user_id, cursor: nil, limit: 20)
    result = TimelineWithZSet.fetch(user_id, cursor: cursor, limit: limit)
    Tweet.where(id: result[:tweet_ids]).includes(:user)
  end
end

セレブリティの昇格・降格バッチ

# lib/tasks/celebrity_classification.rake
namespace :celebrity do
  desc "フォロワー数に基づきユーザー分類を更新"
  task reclassify: :environment do
    # セレブリティ昇格(一般→セレブリティ)
    User.where(is_celebrity: false)
        .where('followers_count >= ?', FanoutService::CELEBRITY_THRESHOLD)
        .find_each do |user|
      user.update!(is_celebrity: true)
      # 既存のPushキャッシュを削除(移行)
      RemovePushCacheJob.perform_later(user.id)
    end
 
    # セレブリティ降格(セレブリティ→一般)
    User.where(is_celebrity: true)
        .where('followers_count < ?', FanoutService::CELEBRITY_THRESHOLD)
        .find_each do |user|
      user.update!(is_celebrity: false)
      # Pushキャッシュを再構築
      RebuildPushCacheJob.perform_later(user.id)
    end
  end
end

Step 7: PostgreSQLのパーティショニング設計

-- tweets テーブル(月単位パーティショニング)
CREATE TABLE tweets (
  id         BIGINT      NOT NULL,  -- Snowflake ID(時系列ソート可能)
  user_id    BIGINT      NOT NULL,
  content    VARCHAR(280) NOT NULL,
  created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
) PARTITION BY RANGE (created_at);
 
-- インデックス設計
CREATE INDEX idx_tweets_user_id_created
  ON tweets (user_id, created_at DESC);  -- プロフィールタイムライン用
 
-- 月単位パーティション
CREATE TABLE tweets_2024_01 PARTITION OF tweets
  FOR VALUES FROM ('2024-01-01') TO ('2024-02-01');
 
-- follows テーブル(シャーディング候補)
CREATE TABLE follows (
  follower_id BIGINT NOT NULL,
  followee_id BIGINT NOT NULL,
  created_at  TIMESTAMPTZ NOT NULL DEFAULT NOW(),
  PRIMARY KEY (follower_id, followee_id)
);
CREATE INDEX idx_follows_followee ON follows(followee_id);
-- followee_idでの逆引き(フォロワー一覧取得)用

Step 8: 検索機能との統合(Amazon OpenSearch)

# app/jobs/tweet_index_job.rb
class TweetIndexJob < ApplicationJob
  queue_as :search_indexing
 
  def perform(tweet_id)
    tweet = Tweet.find(tweet_id)
 
    OPENSEARCH_CLIENT.index(
      index: "tweets-#{tweet.created_at.strftime('%Y-%m')}",
      id:    tweet.id,
      body:  {
        user_id:    tweet.user_id,
        content:    tweet.content,
        created_at: tweet.created_at.iso8601,
        # 検索用のメタデータ
        hashtags: extract_hashtags(tweet.content),
        mentions: extract_mentions(tweet.content)
      }
    )
  end
 
  private
 
  def extract_hashtags(content)
    content.scan(/#\w+/).map(&:downcase)
  end
 
  def extract_mentions(content)
    content.scan(/@\w+/).map { |m| m[1..] }  # @を除去
  end
end
 
# 検索クエリ例
def search_tweets(query, user_id:, cursor: nil)
  OPENSEARCH_CLIENT.search(
    index: 'tweets-*',
    body: {
      query: {
        bool: {
          must: [{ match: { content: query } }],
          filter: [
            # フォロー中のユーザーのみ
            { terms: { user_id: following_ids(user_id) } }
          ]
        }
      },
      sort: [{ created_at: { order: 'desc' } }],
      search_after: cursor ? [cursor] : nil,
      size: 20
    }
  )
end

Step 9: AWSアーキテクチャの詳細

Loading diagram...
Write Path(投稿フロー):
  Client
  → ALB
  → Tweet API(ECS Fargate × 3)
  → Aurora PostgreSQL(永続化)
  → Amazon MSK(Kafka、ファンアウトキュー)
  → Fanout Worker(ECS Fargate × 10台)
  → ElastiCache Redis Cluster(タイムラインキャッシュ)
  → Amazon OpenSearch(全文検索インデックス)
 
Read Path(タイムライン取得フロー):
  Client
  → ALB
  → Timeline API(ECS Fargate × 10)
  → ElastiCache Redis(キャッシュヒット → そのまま返す)
  → Aurora Reader Replica(キャッシュミス時のみ)
  + セレブリティツイートのリアルタイムマージ
 
Amazon MSK(Kafka)採用理由:
  - ファンアウトWorkerがスケールアウトしても確実にメッセージを配信
  - Consumerグループで並列処理
  - at-least-onceのメッセージ配信保証
  - 7日間のメッセージ保持(リプレイ可能)
 
Aurora PostgreSQL:
  Writer 1台 + Reader 3台(Multi-AZ)
  db.r6g.2xlarge(8vCPU, 64GB)
  接続数: PgBouncer経由でコネクションプール管理
 
ElastiCache Redis Cluster:
  cache.r6g.2xlarge × 6ノード(3シャード × 2レプリカ)
  合計メモリ: 312GB
  各ユーザーのタイムライン: 800件 × 8bytes ≒ 6.4KB
  → 100Mユーザー: 640GB(キャッシュに入りきらないのでLRU)

Step 10: 監視設計

# app/middleware/timeline_metrics.rb
class TimelineMetrics
  FANOUT_LATENCY = Prometheus::Client::Histogram.new(
    :fanout_latency_seconds,
    docstring: 'Time to fanout tweet to all followers',
    buckets: [0.1, 0.5, 1.0, 5.0, 10.0, 30.0]
  )
  TIMELINE_CACHE_HIT = Prometheus::Client::Counter.new(
    :timeline_cache_hit_total, docstring: 'Timeline cache hits'
  )
  FANOUT_QUEUE_DEPTH = Prometheus::Client::Gauge.new(
    :fanout_queue_depth, docstring: 'MSK Fanout queue depth'
  )
end
# CloudWatch アラーム設定
アラーム一覧:
  FanoutLatencyHigh:
    メトリクス: p99(fanout_latency_seconds)
    閾値: 30秒超(ファンアウト遅延)
    アクション: Fanout Workerのスケールアウト
 
  MSKConsumerLag:
    メトリクス: kafka.consumer_lag
    閾値: 100,000メッセージ超
    アクション: PagerDuty通知 + Workerスケールアウト
 
  TimelineCacheHitRateLow:
    メトリクス: timeline_cache_hit / total_requests
    閾値: 80%未満
    アクション: Slack通知(キャッシュ設計の見直し検討)
 
  TimelineP99Latency:
    メトリクス: p99(timeline_fetch_latency)
    閾値: 100ms超
    アクション: PagerDuty通知

面接官との厳しい質問と模範回答

面接官: 「Twitterのバードウォッチングアカウントが突然バズって
        フォロワーが1万人を超えた場合、どうなりますか?」

ソウタ: 「その場合はセレブリティ昇格の対象になります。
        ただし即時切り替えは危険なので、バッチ処理で段階的に移行します。
        移行期間中は一時的にPushとPullの両方を使うことになりますが、
        ビジネス的には数分間の遅延は許容できます。」

---

面接官: 「KafkaはAt-least-onceですよね。
        同じツイートが2回ファンアウトされたらどうなりますか?」

ソウタ: 「Redisのタイムラインには同じtweetIDが2回入ります。
        これを防ぐには、LPUSHの前にLPOSで存在チェックする方法がありますが、
        チェックと挿入の間のRace Conditionを避けるにはLuaスクリプトが必要です。
        より実用的には、取得時に重複除去するか、
        tweet_idをSortedSetで管理してZADDの冪等性を利用します。」

---

面接官: 「99.99%の可用性はどう担保しますか?」

ソウタ: 「多層防御で実現します。
        1. ALBによるヘルスチェックと自動フェイルオーバー
        2. Aurora Multi-AZによるDBの自動フェイルオーバー(30秒以内)
        3. ElastiCacheのReader/Writerの自動フェイルオーバー
        4. MSK(Kafka)のマルチAZ構成
        5. ECSサービスのAZをまたいだタスク分散
        単一障害点を排除することで、99.99%(年間52分のダウン)を目指します。」

WARNING

タイムライン設計で「全員にPushで届ける」と言ってしまうと、フォロワー1000万人の有名人が1ツイートするたびに1000万回のRedis書き込みが走る。この「有名人問題」を先に言えると面接官の印象が大きく変わる。


まとめ: 設計の判断ポイント

判断1: Pushにするかフォロワー数で切り替えるか
→ ハイブリッド方式(閾値: フォロワー1万人)

判断2: キャッシュのデータ構造
→ Sorted Set(カーソルベースページネーション対応)
→ データ量が少ない場合はListで十分

判断3: タイムラインキャッシュの保持件数
→ 800件(ほとんどのユーザーはそれ以上遡らない)

判断4: セレブリティツイートのマージタイミング
→ 読み込み時にリアルタイムマージ(鮮度とコストのバランス)

判断5: キャッシュの整合性
→ 結果整合性を許容(数秒の遅延はSNSでは許容範囲)

「今日の一番重要なポイントは何だと思う?」とレイカが聞いた。

「ファンアウトの方式を最初に確認することです。PushかPullか、またはハイブリッドか。それによって設計全体が変わる。そして、有名人問題を先に言えること。」

「その通り。そして、設計の選択には必ず理由がある。その理由をトレードオフとして語れるかが評価されるポイントだ。」

ソウタはノートを閉じた。3ヶ月前、同じ問題で黙り込んだ自分が信じられなかった。

INFO

SNSタイムラインの設計は「スケール問題の総合問題集」。ファンアウト戦略、有名人問題、キャッシュの整合性、ページネーション設計をすべてトレードオフとして語れると、上位候補者として評価される。


次の章では、チャットシステムの設計に挑む。WebSocketとHTTPの違いから始まり、メッセージ配信保証の設計まで。