Проектирование ленты новостей (timeline) -- классический кейс с акцентом на fanout, кэширование и обработку celebrity проблемы.
Шаг 1: Требования
Функциональные требования
- Публикация постов (tweet): текст до 280 символов + медиа
- Home timeline: лента постов от подписок пользователя
- User timeline: все посты конкретного пользователя
- Подписка/отписка (follow/unfollow)
- Поиск по постам
- Trending topics
Нефункциональные требования
- Home timeline latency < 200ms
- Post publication < 5 секунд до появления в лентах подписчиков
- DAU: 300M, Peak tweet rate: ~10K/sec
- 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);
}
}
Шаг 5: Trending Topics
<?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 |
Возможные вопросы интервьюера
-
Как обрабатывать celebrity с 50M подписчиков?
- Гибридный подход: push для обычных, pull для celebrity
- Порог ~100K подписчиков
- Merge на чтении
-
Как обеспечить порядок в timeline?
- Snowflake ID содержит timestamp
- Redis sorted set с timestamp score
- Merge sort при гибридном подходе
-
Как удалить tweet из всех timeline?
- Мягкое удаление: пометить deleted в БД
- Lazy cleanup: фильтровать при чтении
- Background job для очистки кэшей
-
Как обрабатывать retweets?
- Retweet = новый tweet с reference на original
- Фанаут как обычный tweet
- Денормализация данных original tweet
-
Cache stampede при celebrity tweet?
- Staggered fanout (не все сразу)
- Rate limiting на fanout workers
- Priority queue для online пользователей