HardКейс5 min

Twitter/X Feed

Проектирование ленты новостей: fanout on write vs read, timeline generation, trending topics

Проектирование ленты новостей (timeline) -- классический кейс с акцентом на fanout, кэширование и обработку celebrity проблемы.

Шаг 1: Требования

Функциональные требования

  1. Публикация постов (tweet): текст до 280 символов + медиа
  2. Home timeline: лента постов от подписок пользователя
  3. User timeline: все посты конкретного пользователя
  4. Подписка/отписка (follow/unfollow)
  5. Поиск по постам
  6. Trending topics

Нефункциональные требования

  1. Home timeline latency < 200ms
  2. Post publication < 5 секунд до появления в лентах подписчиков
  3. DAU: 300M, Peak tweet rate: ~10K/sec
  4. Eventual consistency допустима

Шаг 2: Оценка нагрузки

Метрика Значение
DAU 300M
Tweets в день 500M
Среднее подписок 200
Timeline reads в день 30B
Read QPS ~350K
Write QPS ~6K
Read/Write ratio ~60:1

Шаг 3: High-Level архитектура

┌──────────┐     ┌───────────────┐     ┌────────────────────────────────┐
│  Client  │────>│  API Gateway  │────>│  Tweet Service                 │
│          │     │               │     │  (publish, get)                │
└──────────┘     └───────────────┘     └───────────┬────────────────────┘
                                                   │
                                        ┌──────────▼──────────┐
                                        │  Fanout Service     │
                                        │  (write to caches)  │
                                        └──────────┬──────────┘
                                                   │
                    ┌──────────────────────────────┼────────────────┐
                    │                              │                │
             ┌──────▼──────┐           ┌───────────▼──────┐  ┌─────▼──────────┐
             │  Timeline   │           │  Tweet Store     │  │  Social Graph  │
             │  Cache      │           │  (PostgreSQL)    │  │  (followers)   │
             │  (Redis)    │           └──────────────────┘  └────────────────┘
             └─────────────┘

Шаг 4: Fanout on Write vs Fanout on Read

Fanout on Write (push модель)

<?php

declare(strict_types=1);

/**
 * При публикации tweet -- заранее записываем его в timeline cache
 * каждого подписчика.
 */
final class FanoutOnWriteService
{
    public function __construct(
        private readonly \Redis $redis,
        private readonly FollowerRepository $followers,
        private readonly TweetRepository $tweets,
        private readonly QueuePublisher $queue,
    ) {}

    public function publishTweet(string $userId, string $content): string
    {
        // 1. Save tweet
        $tweetId = $this->tweets->create($userId, $content);

        // 2. Get all followers
        $followerIds = $this->followers->getFollowerIds($userId);

        // 3. Fanout to each follower's timeline cache
        foreach (array_chunk($followerIds, 1000) as $batch) {
            $this->queue->publish('fanout.batch', [
                'tweet_id' => $tweetId,
                'user_id' => $userId,
                'follower_ids' => $batch,
                'timestamp' => microtime(true),
            ]);
        }

        // 4. Add to own timeline
        $this->addToTimeline($userId, $tweetId);

        return $tweetId;
    }

    /**
     * Worker: process fanout batch
     */
    public function processFanoutBatch(array $job): void
    {
        $tweetEntry = json_encode([
            'tweet_id' => $job['tweet_id'],
            'user_id' => $job['user_id'],
            'timestamp' => $job['timestamp'],
        ]);

        $pipeline = $this->redis->pipeline();

        foreach ($job['follower_ids'] as $followerId) {
            $key = "timeline:{$followerId}";
            $pipeline->zAdd($key, $job['timestamp'], $tweetEntry);

            // Keep only latest 800 tweets in cache
            $pipeline->zRemRangeByRank($key, 0, -801);
        }

        $pipeline->exec();
    }

    /**
     * Read home timeline (fast -- just read from cache)
     */
    public function getHomeTimeline(
        string $userId,
        int $limit = 20,
        ?float $beforeTimestamp = null,
    ): array {
        $key = "timeline:{$userId}";
        $max = $beforeTimestamp ?? '+inf';

        $entries = $this->redis->zRevRangeByScore(
            $key,
            (string) $max,
            '-inf',
            ['limit' => [0, $limit]],
        );

        if (empty($entries)) {
            // Cache miss -- build from DB
            return $this->buildTimelineFromDb($userId, $limit);
        }

        // Hydrate tweet data
        $tweetIds = array_map(
            fn (string $entry) => json_decode($entry, true)['tweet_id'],
            $entries,
        );

        return $this->tweets->getByIds($tweetIds);
    }

    private function addToTimeline(string $userId, string $tweetId): void
    {
        $key = "timeline:{$userId}";
        $entry = json_encode([
            'tweet_id' => $tweetId,
            'user_id' => $userId,
            'timestamp' => microtime(true),
        ]);

        $this->redis->zAdd($key, microtime(true), $entry);
        $this->redis->zRemRangeByRank($key, 0, -801);
    }
}

Плюсы: Чтение timeline мгновенное (O(1) из Redis). Минусы: Celebrities с миллионами подписчиков создают огромный fanout.

Fanout on Read (pull модель)

<?php

declare(strict_types=1);

/**
 * При чтении timeline -- собираем tweets от всех подписок на лету.
 */
final class FanoutOnReadService
{
    public function __construct(
        private readonly TweetRepository $tweets,
        private readonly FollowingRepository $following,
        private readonly \Redis $redis,
    ) {}

    public function getHomeTimeline(
        string $userId,
        int $limit = 20,
        ?string $cursor = null,
    ): array {
        // 1. Get list of people user follows
        $followingIds = $this->following->getFollowingIds($userId);

        // 2. Get recent tweets from each (with cache)
        $allTweets = [];

        foreach ($followingIds as $followedId) {
            $tweets = $this->getUserRecentTweets($followedId, 50);
            $allTweets = array_merge($allTweets, $tweets);
        }

        // 3. Sort by timestamp descending
        usort($allTweets, fn ($a, $b) => $b['timestamp'] <=> $a['timestamp']);

        // 4. Apply cursor-based pagination
        if ($cursor !== null) {
            $allTweets = array_filter(
                $allTweets,
                fn ($t) => $t['timestamp'] < $cursor,
            );
        }

        return array_slice($allTweets, 0, $limit);
    }

    private function getUserRecentTweets(string $userId, int $limit): array
    {
        $cacheKey = "user_tweets:{$userId}";
        $cached = $this->redis->get($cacheKey);

        if ($cached !== false) {
            return json_decode($cached, true);
        }

        $tweets = $this->tweets->getByUser($userId, $limit);
        $this->redis->setex($cacheKey, 300, json_encode($tweets)); // 5 min cache

        return $tweets;
    }
}

Плюсы: Нет fanout overhead при публикации. Минусы: Медленное чтение (нужно собирать tweets от всех подписок).

Гибридный подход (рекомендуемый)

<?php

declare(strict_types=1);

final class HybridTimelineService
{
    private const CELEBRITY_THRESHOLD = 100_000; // followers

    public function __construct(
        private readonly FanoutOnWriteService $pushService,
        private readonly FanoutOnReadService $pullService,
        private readonly FollowerRepository $followers,
        private readonly TweetRepository $tweets,
        private readonly \Redis $redis,
    ) {}

    public function publishTweet(string $userId, string $content): string
    {
        $tweetId = $this->tweets->create($userId, $content);
        $followerCount = $this->followers->getFollowerCount($userId);

        if ($followerCount < self::CELEBRITY_THRESHOLD) {
            // Regular user: fanout on write
            $this->pushService->fanoutToFollowers($tweetId, $userId);
        } else {
            // Celebrity: DON'T fanout, will be merged on read
            $this->redis->sAdd('celebrities', $userId);
        }

        return $tweetId;
    }

    public function getHomeTimeline(string $userId, int $limit = 20): array
    {
        // 1. Get pre-computed timeline (from fanout on write)
        $cachedTimeline = $this->pushService->getHomeTimeline($userId, $limit * 2);

        // 2. Get tweets from celebrities user follows
        $followingIds = $this->followers->getFollowingIds($userId);
        $celebIds = $this->redis->sInter('celebrities', ...$followingIds);

        $celebTweets = [];
        foreach ($celebIds as $celebId) {
            $tweets = $this->pullService->getUserRecentTweets($celebId, 20);
            $celebTweets = array_merge($celebTweets, $tweets);
        }

        // 3. Merge and sort
        $merged = array_merge($cachedTimeline, $celebTweets);
        usort($merged, fn ($a, $b) => $b['timestamp'] <=> $a['timestamp']);

        // 4. Deduplicate
        $seen = [];
        $result = [];
        foreach ($merged as $tweet) {
            if (!isset($seen[$tweet['id']])) {
                $seen[$tweet['id']] = true;
                $result[] = $tweet;
            }
        }

        return array_slice($result, 0, $limit);
    }
}
<?php

declare(strict_types=1);

final class TrendingService
{
    private const WINDOW_MINUTES = 60;
    private const TOP_N = 10;

    public function __construct(
        private readonly \Redis $redis,
    ) {}

    /**
     * Track hashtag usage (called on each tweet)
     */
    public function trackHashtags(array $hashtags): void
    {
        $currentWindow = (int) (time() / 60); // per-minute window

        foreach ($hashtags as $hashtag) {
            $tag = mb_strtolower($hashtag);

            // Increment in current minute window
            $key = "trending:{$currentWindow}";
            $this->redis->zIncrBy($key, 1.0, $tag);
            $this->redis->expire($key, self::WINDOW_MINUTES * 60 + 60);
        }
    }

    /**
     * Get current trending topics
     */
    public function getTrending(int $limit = 10): array
    {
        $currentWindow = (int) (time() / 60);
        $keys = [];

        // Aggregate last N minutes
        for ($i = 0; $i < self::WINDOW_MINUTES; $i++) {
            $keys[] = "trending:" . ($currentWindow - $i);
        }

        // Union all windows into a temporary key
        $resultKey = 'trending:current';
        $weights = [];

        // More recent windows get higher weight
        for ($i = 0; $i < count($keys); $i++) {
            $weights[] = 1.0 / ($i + 1); // Decay factor
        }

        $this->redis->zUnionStore($resultKey, $keys, $weights, 'SUM');
        $this->redis->expire($resultKey, 60);

        // Get top N
        $trending = $this->redis->zRevRange($resultKey, 0, $limit - 1, true);

        return array_map(
            fn (string $tag, float $score) => [
                'hashtag' => $tag,
                'score' => round($score),
            ],
            array_keys($trending),
            array_values($trending),
        );
    }
}

Шаг 6: Масштабирование

Компонент Стратегия
Timeline cache Redis Cluster (sharded by user_id)
Tweet store PostgreSQL sharded by tweet_id
Social graph Graph DB или PostgreSQL sharded
Search Elasticsearch
Fanout workers Auto-scaling (burst при celebrity posts)
Trending Redis + Flink для real-time aggregation

Возможные вопросы интервьюера

  1. Как обрабатывать celebrity с 50M подписчиков?

    • Гибридный подход: push для обычных, pull для celebrity
    • Порог ~100K подписчиков
    • Merge на чтении
  2. Как обеспечить порядок в timeline?

    • Snowflake ID содержит timestamp
    • Redis sorted set с timestamp score
    • Merge sort при гибридном подходе
  3. Как удалить tweet из всех timeline?

    • Мягкое удаление: пометить deleted в БД
    • Lazy cleanup: фильтровать при чтении
    • Background job для очистки кэшей
  4. Как обрабатывать retweets?

    • Retweet = новый tweet с reference на original
    • Фанаут как обычный tweet
    • Денормализация данных original tweet
  5. Cache stampede при celebrity tweet?

    • Staggered fanout (не все сразу)
    • Rate limiting на fanout workers
    • Priority queue для online пользователей