Cache Celebrity Problem: почему 100 Redis-нод не спасут вас от одного hot key — и что заложить в архитектуру до первого крупного клиента

Cache Celebrity Problem: почему 100 Redis-нод не спасут вас от одного hot key — и что заложить в архитектуру до первого крупного клиента

Можно идеально разложить миллиард ключей по 100 Redis-нодам и всё равно получить деградацию из-за одного merchant:123. Если один клиент внезапно начинает генерировать 500 000 RPS, средняя загрузка кластера может выглядеть прекрасно ровно до того момента, когда один shard упрётся в CPU и сеть.

В распределённых системах мы привыкли лечить нагрузку горизонтальным масштабированием: добавили ноды, разбили данные на shard'ы, распределили ключи через hash или consistent hashing — живём дальше.

Проблема в том, что эта модель хорошо работает, пока нагрузка распределена более-менее равномерно.

В реальной системе она почти никогда не равномерна.

Один merchant проводит платежей на два порядка больше остальных. Один SKU внезапно становится вирусным. Один пост читают миллионы пользователей. Один tenant вырастает так, что начинает потреблять заметную долю всего кластера.

Так появляется Cache Celebrity Problem, она же частный случай hot key / hot partition problem.

Почему вообще «celebrity»

Проще всего объяснить на социальной сети.

Допустим, последние посты пользователя лежат в кеше так:

user_posts:{user_id}​

У обычного пользователя этот ключ читают несколько раз в минуту.

Теперь Cristiano Ronaldo публикует пост. Миллионы пользователей начинают читать:

user_posts:ronaldo

С точки зрения данных ничего особенного не произошло. Это по-прежнему один ключ.

С точки зрения нагрузки этот ключ отличается от остальных на несколько порядков.

Отсюда и название: небольшое количество «знаменитостей» получает непропорционально большую долю трафика.

Причём celebrity — это необязательно человек. В платёжной системе им легко становится крупный merchant или банковский BIN. В e-commerce — популярный SKU. В SaaS — жирный tenant. В инфраструктуре — конфигурационный ключ, который читает весь fleet.

Главное свойство одно:

traffic(key) >> average traffic(key)

То есть проблема не в объёме данных. Проблема в skew распределения нагрузки.

Почему средняя загрузка врёт

Представим Redis Cluster из 100 нод.

Есть миллиард ключей, которые отлично распределились по shard'ам:

hash(key) -> shard

На каждый shard приходится примерно одинаковый объём данных. Красиво.

Потом появляется:

merchant:123

И получает:

500 000 RPS​

Hash всё равно отправляет ключ на одну ноду:

hash(merchant:123) -> shard-42

Получаем примерно это:

500k RPS | v shard-42 / \ CPU 100% NET 100%

Остальные 99 shard'ов в этот момент могут быть загружены на 5%.

Средний CPU кластера будет выглядеть почти издевательски хорошо:

average CPU = 6%

А система уже деградирует.

Это одна из самых неприятных частей Celebrity Problem: агрегированные метрики успокаивают ровно тогда, когда надо начинать нервничать.

Если смотреть только на average CPU, average latency и суммарный RPS, hotspot легко пропустить. Нужны метрики уровня конкретного shard'а и конкретной нагрузки:

max shard CPU max shard RPS max shard network top-K keys by RPS top-K tenants by RPS

Особенно полезны отношения вроде:

max_shard_rps / avg_shard_rps

и:

top_1_tenant_rps / total_rps

Если средний shard держит 10k RPS, а максимальный уже 80k RPS, то skew factor = 8. Средняя температура по больнице здесь ничего не объясняет.

Почему consistent hashing не спасает

У consistent hashing другая работа.

Он хорошо распределяет разные ключи между нодами:

A -> shard-1 B -> shard-7 C -> shard-3 D -> shard-9​

Но один ключ получает одно положение:

HOT_KEY -> shard-7

Если половина всего трафика системы приходит в этот ключ, половина трафика окажется на одном shard'е.

Отсюда принцип, который полезно держать в голове при любом шардировании:

балансировка данных != балансировка нагрузки

Архитектору куда важнее знать распределение RPS по ключам и tenant'ам, чем радоваться одинаковому количеству ключей на shard.

Поэтому и очевидные решения часто не работают.

«Redis перегружен — давайте добавим нод» лечит общую нагрузку. Если было 10 shard'ов, стало 100, hot key просто переедет с shard-7 на shard-73. Инфраструктура станет дороже в десять раз, а bottleneck останется на одной машине.

Можно поставить shard помощнее. Иногда это нормальная временная мера. Но 500k RPS -> one machine остаётся той же архитектурой. Завтра будет 1M RPS — и потолок вернётся.

Ещё хуже бывает с TTL. Интуитивно хочется уменьшить его ради свежести. Но если горячий ключ получает 100k RPS и протухает каждую секунду, один cache miss легко превращается в толпу одинаковых запросов к downstream.

Было:

hot cache node

Стало:

dead database

Это уже cache stampede / thundering herd.

Что действительно помогает hot key

Универсальной кнопки нет. Нормальная защита обычно строится слоями.

L1 cache: самый дешёвый запрос — тот, которого не было

Если значение допускает хотя бы секунду stale-данных, локальный кеш внутри application instance может снять несколько порядков запросов с Redis.

Вместо:

Application | v Redis

получаем:

Application | v L1 cache | v Redis

Например:

local cache TTL = 1 sec Redis TTL = 60 sec

Есть 100 application instances и 100k RPS на один горячий объект.

Без L1 это до 100k Redis reads/sec.

С L1 при достаточно горячих данных мы можем получить порядок 100 Redis reads/sec: каждый instance обновляет локальную копию примерно раз в секунду.

Цена понятна: данные между instances некоторое время расходятся, а холодный старт fleet способен ударить вниз по стеку. Но если бизнес допускает секундную stale-версию, это очень дешёвый способ убрать сетевой трафик из самой горячей части path.

Репликация hot key: одному logical key нужны несколько физических копий

Если один ключ уже не помещается по RPS на одну машину, нет смысла продолжать хранить одну физическую копию.

Вместо:

celebrity:42

делаем:

celebrity:42:0 celebrity:42:1 celebrity:42:2 ... celebrity:42:N

На чтении выбираем одну из реплик:

replica := rand.Intn(N) key := fmt.Sprintf( "celebrity:%d:%d", userID, replica, )

Разные physical keys попадут на разные shard'ы, и один logical key начнёт масштабироваться горизонтально.

Это, пожалуй, главный сдвиг мышления в Celebrity Problem:

если один logical key слишком горячий для одной машины, ему нужны несколько physical representations.

Мы больше не пытаемся масштабировать «кластер вообще». Мы масштабируем конкретную celebrity.

Request coalescing: 10 000 cache miss должны стать одним запросом

Репликация спасает steady-state чтения, но не cache miss.

Допустим, 10 000 запросов одновременно хотят объект, которого сейчас нет в кеше.

Плохой вариант:

10 000 cache misses | v 10 000 DB requests

Нормальный:

10 000 cache misses | v singleflight | v 1 DB request

Первый запрос загружает данные. Остальные ждут его результат.

В Go на одном instance это может быть обычный:

singleflight.Group

Для distributed-варианта понадобится coordination между instances или request collapsing ниже по стеку.

Здесь важен сам принцип: popular cache miss опаснее popular cache hit. Пока значение есть в кеше, мы мучаем cache layer. Когда оно исчезло, celebrity способна мгновенно протащить весь свой RPS в базу.

Soft TTL + hard TTL: протухание не обязано быть бинарным

Классическая модель TTL выглядит так:

fresh -> deleted

Для hot key это довольно нервная схема.

Полезнее иметь промежуточное состояние:

fresh -> stale -> expired

Например:

Client A -> stale -> refresh Client B -> stale -> return old value Client C -> stale -> return old value

Это stale-while-revalidate / lazy revalidation.

Тот же принцип хорошо работает во время проблем downstream. Иногда старые данные на пару минут дешевле, чем идеально свежий ответ ценой умершей базы.

Pre-warming: если знаем о всплеске заранее, зачем ждать miss

Некоторые celebrity предсказуемы.

Пользователь с десятками миллионов подписчиков публикует пост. Большой merchant запускает распродажу. Популярный товар появляется в промо.

В таких случаях можно прогреть cache replicas заранее:

celebrity publishes post | +--> DB | +--> cache replica 1 +--> cache replica 2 +--> cache replica 3 +--> cache replica N

Для news feed это особенно хорошо ложится на гибридную схему. Обычных пользователей можно обслуживать через push/fan-out, а celebrity перевести на pull и держать отдельно закешированной.

Не надо пушить запись в десятки миллионов timelines. Достаточно заранее положить её в ограниченное число cache replicas.

Когда celebrity — это жирный tenant

В платёжных и SaaS-системах интереснее другой вариант.

Допустим, есть 100 000 клиентов. Большинство делает порядка 10 RPS.

Но несколько выглядят так:

merchant-A = 10 000 RPS merchant-B = 20 000 RPS merchant-C = 50 000 RPS

А routing устроен классически:

shard = hash(merchant_id) % N

В какой-то момент крупные merchants случайно оказываются вместе:

shard-7: merchant-A 10k merchant-B 20k merchant-C 50k others 5k TOTAL = 85k RPS

Остальные shard'ы получают около 10k RPS.

Это та же Celebrity Problem. Только celebrity теперь — жирный клиент.

И вот здесь hash(tenant_id) -> shard начинает мешать.

Пока система маленькая, жёсткая hash-схема очень удобна. Никакого metadata storage, routing map и control plane. Взяли tenantID, посчитали hash, нашли shard.

Проблема появляется в тот день, когда нужно сказать:

merchant-A теперь живёт на shard-42​

А hash-функция отвечает: «Нет».

Поэтому в крупной системе полезен слой indirection:

tenantID | v Shard Map | v shardID

Для большинства tenants mapping всё ещё можно получать обычным hash:

default -> hash(tenantID)

Для исключений появляется override:

overrides[tenantID] -> dedicated shard

В коде реализауия примитивная:

func shardForTenant(id TenantID) ShardID { if shard, ok := overrides[id]; ok { return shard } return consistentHash(id) }

Зато архитектурно меняется очень многое. Placement становится управляемым.

Как вынести крупного клиента без большого downtime

Если речь про cache, миграция относительно простая.

Сначала поднимаем target shard. Потом прогреваем на нём самые горячие ключи. Если переключить жирного клиента на пустой cache, migration сама создаст cache stampede.

После прогрева переключаем routing:

было: merchant-A -> shard-7 стало: merchant-A -> shard-42

Mapping лучше версионировать и быстро распространять между application instances:

tenant_shards: merchant-A: shard: 42 version: 172

Старый cache я бы не удалял сразу. Пусть естественно умирает по TTL. Тогда rollback — это смена mapping обратно на 42 -> 7, пока старая копия ещё жива.

С persistent database shard всё заметно сложнее: backfill, CDC/replication, сверка состояния, при необходимости shadow reads, переключение routing, rollback window и только потом удаление старых данных.

Но главный вывод тот же:

routing не должен быть намертво зашит в hash-функцию.

Нужен control plane, который умеет перемещать конкретную нагрузку.

Иначе после появления первого действительно крупного клиента выяснится странная вещь: кластер масштабировать мы умеем, а одного клиента — нет.

Dedicated shard тоже закончится

Допустим, жирного merchant вынесли на отдельный shard.

Работает.

Через год он вырос ещё в десять раз и стал больше capacity одной машины.

Мы снова в исходной точке.

Тогда появляется следующий уровень:

normal tenants | v shared shards large tenant | v dedicated shard very large tenant | v tenant internal sharding

То есть сначала система шардируется по клиентам. Потом конкретного клиента приходится шардировать уже внутри него самого:

hash(merchant_id, account_id) hash(merchant_id, payment_id) hash(merchant_id, order_id)

Celebrity сначала перестаёт делить shard с соседями, а потом сама становится небольшим распределённым кластером.

И это нормальная эволюция. Плохой сценарий — когда архитектура считает, что любой tenant навсегда меньше одного shard.

Что я бы закладывал в систему заранее

Главная ошибка в таких системах — предполагать, что равномерное распределение данных автоматически даст равномерную нагрузку.

Не даст.

Один logical entity способен генерировать трафик, сопоставимый со всем остальным кластером. Поэтому система должна уметь обнаружить hotspot, изолировать его, размазать его нагрузку и при необходимости изменить placement.

Практические правила у меня здесь простые.

  • Consistent hashing балансирует ключи, а не RPS.
  • Scale-out кластера не лечит single hot key.
  • Popular cache miss опаснее popular cache hit, поэтому нужны singleflight, stale-while-revalidate и pre-warming.
  • Крупные tenants надо уметь выносить на dedicated shards.
  • Routing через чистую hash-функцию удобен ровно до первого серьёзного resharding. Дальше нужен управляемый слой tenant -> shard.
  • Если один tenant перерастает capacity отдельного shard, его самого придётся шардировать.

В конечном счёте задача не в том, чтобы красиво и равномерно разложить данные. Задача — управляемо разложить нагрузку.

Именно эта разница обычно отделяет систему, которая отлично выглядит на архитектурной диаграмме, от системы, которая нормально переживает первого по-настоящему большого клиента.

Если вам интересны такие разборы про highload, distributed systems, Go и финтех — подписывайтесь на мой Telegram-канал.

Селькин Андрей
TGM / Go-разработчик из Fintech. Инженерные заметки: https://t.me/andrei_selkin_outbox
1