Кэширование и масштабирование: уровни, инвалидация, шардирование
Есть два способа сделать систему быстрее: делать меньше работы или делать работу на большем количестве машин. Кэширование — это первое. Масштабирование — второе. Их объединяют в одну статью не по лени, а потому что это две стороны одного разговора: почти всегда правильный порядок действий — «сначала перестань считать одно и то же по десять раз, потом добавляй железо». Инженер, который начинает с шардирования базы, обычно через полгода обнаруживает, что 80% запросов к шардам — это один и тот же справочник тарифов, который можно было держать в памяти процесса.
При этом кэш — самый коварный инструмент в арсенале. Он не ломается громко. Он тихо отдаёт данные вчерашнего дня, и вы узнаёте об этом от службы поддержки через две недели. Фил Карлтон сформулировал это раз и навсегда: «В компьютерных науках только две сложные вещи — инвалидация кэша и придумывание имён». Эта статья — в основном про первую.
Предполагается знакомство со статьями «Микросервисы: границы, коммуникация, данные, эксплуатация» и «Событийная архитектура»: понятия «владелец данных», «событие», «партиция» дальше используются без объяснений.
1. Прежде чем масштабировать: три закона, которые нельзя обойти
Масштабирование обсуждают эмоционально («добавим подов»), а оно подчиняется арифметике. Три соотношения дают почти всё, что нужно для честного разговора.
Закон Литтла
Для любой стабильной системы с очередью:
L = λ · W
где L — среднее число запросов внутри системы, λ — частота поступления (запросов в секунду), W — среднее время в системе. Закон не требует никаких предположений о распределениях — он верен всегда, когда система стабильна (J. D. C. Little, 1961; практическое изложение — глава 13 «Systems Performance» Брендана Грегга).
Практическая польза мгновенная. Сервис держит 2000 rps при среднем времени ответа 50 мс. Значит внутри системы в среднем 2000 × 0.05 = 100 запросов. Если ваш пул потоков — 40, вы уже стоите в очереди, и время ответа не 50 мс, а больше. Отсюда же формула размера пула соединений к БД: если запрос к БД занимает 5 мс и вам нужно 3000 запросов к БД в секунду, минимум 3000 × 0.005 = 15 соединений. Не 200, как обычно ставят «на всякий случай» — лишние соединения только увеличивают конкуренцию за защёлки в базе.
Закон Амдала
Если доля p работы распараллеливаема, а (1−p) — строго последовательна, ускорение на N процессорах ограничено:
S(N) = 1 / ((1 − p) + p/N)
При p = 0.95 предел ускорения — 20×, сколько бы машин вы ни добавили. В архитектуре «последовательная часть» — это обычно единственная строка в БД, по которой берётся блокировка, единственный лидер, глобальный счётчик или общий Redis-инстанс.
Универсальный закон масштабируемости (USL)
Нил Гюнтер добавил к Амдалу второй член — стоимость когерентности, то есть согласования между узлами (Gunther, «Guerrilla Capacity Planning»):
C(N) = N / (1 + α(N − 1) + βN(N − 1))
α — доля последовательного (contention), β — стоимость когерентности (crosstalk). Ключевое отличие: при β > 0 пропускная способность не выходит на плато, а падает после некоторого N. Это то самое явление, когда вы добавили 20 подов, и стало хуже. Причина обычно тривиальна: каждый под открывает соединения к БД, каждый пишет в общий кэш, каждый участвует в gossip — и квадратичный член побеждает.
Вывод для практика: прежде чем добавлять узлы, спросите, что у них общее. Общая база, общий кэш, общий лидер выборов — это α и β. Кэш, кстати, помогает именно тем, что убирает обращения к общему ресурсу, то есть уменьшает α.
2. Куб масштабирования: три независимые оси
Мартин Эббот и Майкл Фишер в «The Art of Scalability» предложили модель, которая до сих пор остаётся лучшей картой (изложение у Криса Ричардсона):
Оси ортогональны, и порядок применения почти всегда один и тот же: X → Y → Z. Клонирование бесплатно (если у вас stateless), функциональная декомпозиция стоит организационных усилий (см. «Микросервисы»), а Z-разделение необратимо усложняет модель данных, поэтому его откладывают до последнего.
Отдельно стоит вертикальное масштабирование — просто взять машину побольше. Его недооценивают из снобизма. Современный сервер — 128 ядер и 2 ТБ RAM; база, которая целиком помещается в память, ведёт себя качественно иначе. Если вы можете купить ещё три года на одной большой машине за 2000 $/мес вместо шести месяцев работы команды над шардированием — это выгодная сделка. Подробнее о недооценённости простых решений — в статье «Монолит и модульный монолит».
3. Уровни кэширования
Кэш — не одно место, а лестница. Каждая ступень быстрее следующей примерно на порядок и «протухшее» её примерно на столько же.
Путь одного HTTP-запроса через все уровни:
Cache-Control: max-age} B -->|hit| R1([Ответ, 0 мс, 0 байт по сети]) B -->|miss| CDN{Кэш CDN / edge} CDN -->|hit| R2([Ответ, ~20 мс]) CDN -->|stale + SWR| R2 CDN -->|miss| GW[API Gateway / reverse proxy] GW --> APP[Инстанс приложения] APP --> L1{Локальный кэш
в памяти процесса} L1 -->|hit| R3([Ответ, ~1 мс]) L1 -->|miss| L2{Распределённый кэш
Redis / Memcached} L2 -->|hit| R4([Ответ, ~2 мс]) L2 -->|miss| DB[(База данных)] DB --> MAT{Материализованное
представление?} MAT -->|да| Q1[Готовая агрегация] MAT -->|нет| Q2[Полный запрос + JOIN] Q1 --> FILL[Запись в L2 и L1] Q2 --> FILL FILL --> R5([Ответ, 10–100 мс]) style R1 fill:#5aa469,stroke:#5aa469,color:#fff style R2 fill:#5aa469,stroke:#5aa469,color:#fff style R3 fill:#7cb342,stroke:#7cb342,color:#fff style R4 fill:#c98a2b,stroke:#c98a2b,color:#fff style R5 fill:#d9534f,stroke:#d9534f,color:#fff
Разберём ступени содержательно.
Кэш браузера и HTTP-кэш. Самый дешёвый уровень: ответ вообще не покидает устройство. Управляется заголовками Cache-Control, ETag, Last-Modified. Правило: статику с хешем в имени (app.a3f9c1.js) отдавайте с max-age=31536000, immutable, HTML — с no-cache (что означает не «не кэшировать», а «кэшируй, но всегда перепроверяй»).
CDN / edge. Кэш на географически близких к пользователю узлах. Убирает трансконтинентальный RTT, который иначе неустраним. Ключевая современная возможность — stale-while-revalidate: пользователь получает чуть устаревший ответ мгновенно, а обновление происходит в фоне. Подробнее об edge-архитектурах — в статье «Serverless и edge-архитектуры».
Локальный кэш в процессе (L1). Caffeine в JVM, lru_cache в Python, sync.Map/Ristretto в Go. Наносекундный доступ, ноль сети. Две проблемы: (1) у каждой реплики свой кэш, значит N реплик = N разных версий данных и N промахов при холодном старте; (2) инвалидация требует широковещательного сообщения всем репликам. Годится для того, что меняется редко: справочники, фиче-флаги, конфигурация, скомпилированные политики доступа.
Распределённый кэш (L2). Redis, Memcached. Общий для всех реплик, значит одна версия истины и одна инвалидация. Стоит сетевой round-trip. Это рабочая лошадь: сессии, результаты дорогих запросов, счётчики, лимитеры.
Уровень базы. Материализованные представления, denormalized-таблицы, покрывающие индексы — тоже кэш, просто с транзакционной согласованностью. Часто это лучший ответ: REFRESH MATERIALIZED VIEW CONCURRENTLY решает задачу без нового компонента в инфраструктуре. Крайняя форма этого подхода — отдельная модель чтения, см. «CQRS и Event Sourcing».
Двухуровневый кэш: когда он оправдан
Комбинация L1+L2 даёт лучшее из двух миров, но добавляет второй горизонт устаревания. Оправдана, когда есть выраженные горячие ключи: 1% ключей даёт 50% трафика, и этот 1% помещается в память процесса.
import time
from typing import Any, Callable, Optional
class TwoTierCache:
"""L1 — локальный словарь с коротким TTL, L2 — Redis с длинным.
Инвариант: L1 TTL << L2 TTL. Локальный кэш может отставать
не больше чем на свой TTL — это осознанный бюджет рассогласования.
"""
def __init__(self, redis, l1_ttl: float = 5.0, l2_ttl: int = 300):
self._redis = redis
self._l1: dict[str, tuple[Any, float]] = {}
self._l1_ttl = l1_ttl
self._l2_ttl = l2_ttl
def get(self, key: str, loader: Callable[[], Any]) -> Any:
now = time.monotonic()
# L1: попадание — сотни наносекунд
entry = self._l1.get(key)
if entry is not None and entry[1] > now:
return entry[0]
# L2: попадание — примерно 0,5 мс
raw = self._redis.get(key)
if raw is not None:
value = deserialize(raw)
self._l1[key] = (value, now + self._l1_ttl)
return value
# Промах на обоих уровнях — идём к источнику истины
value = loader()
self._redis.set(key, serialize(value), ex=self._l2_ttl)
self._l1[key] = (value, now + self._l1_ttl)
return value
Стоимость: O(1) на чтение, память O(размер L1) на каждую реплику. Опасность: _l1 без ограничения размера — это утечка памяти. В проде берите библиотеку с вытеснением (Caffeine, cachetools.TTLCache), а не голый словарь.
4. Стратегии чтения и записи
Пять канонических паттернов. Разница между ними — в том, кто отвечает за наполнение кэша и когда данные попадают в хранилище.
| Паттерн | Кто читает из БД | Когда пишется в БД | Риск потери | Где применяют |
|---|---|---|---|---|
| Cache-aside (lazy loading) | приложение | приложение, напрямую | нет | 90% случаев, дефолт |
| Read-through | сам кэш (через провайдера) | — | нет | Caffeine LoadingCache, EHCache |
| Write-through | — | кэш синхронно пишет в БД | нет | когда нужна свежесть кэша |
| Write-behind (write-back) | — | кэш пишет асинхронно, батчами | да, при падении | счётчики, метрики, лайки |
| Refresh-ahead | кэш заранее, до истечения TTL | — | нет | предсказуемо горячие ключи |
Cache-aside — умолчание, потому что он честен: кэш ничего не знает о БД, БД ничего не знает о кэше, приложение управляет обоими. Минус — три места в коде, где легко ошибиться, и «первый запрос всегда медленный».
Write-behind заслуживает отдельного предупреждения. Он превращает кэш в источник истины на время окна записи. Если Redis без AOF-персистентности упадёт, вы потеряете подтверждённые пользователю операции. Применять только там, где потеря нескольких секунд данных допустима: счётчики просмотров — да, списание денег — категорически нет.
Refresh-ahead — недооценённый паттерн. Вместо того чтобы ждать истечения TTL и получить промах на горячем ключе, обновляем запись заранее, когда до истечения осталось меньше порога. Требует знания, какие ключи горячие, но убирает целый класс всплесков латентности.
Жизненный цикл записи в кэше удобно держать в голове как автомат:
5. Инвалидация: почему это действительно трудно
Инвалидация трудна потому, что это распределённая задача согласования между двумя системами, которые не участвуют в общей транзакции. Есть три подхода, и все три несовершенны.
5.1 TTL — истечение по времени
Самый простой и самый надёжный. Ключ живёт T секунд, потом исчезает. Данные могут отставать максимум на T.
Достоинство: система самовосстанавливается. Любая ошибка инвалидации, любой пропущенный сигнал, любое рассогласование — всё чинится само за T секунд. Это огромная ценность, которую недооценивают.
Недостаток: вы всегда отдаёте потенциально устаревшие данные и всегда делаете лишнюю работу для данных, которые не менялись.
Практический приём — джиттер TTL. Если 10 000 ключей записаны одним прогревом с TTL=300, они истекут одновременно, и вы получите синхронный шторм промахов (cache avalanche). Лечится тривиально:
import random
def ttl_with_jitter(base: int, spread: float = 0.2) -> int:
"""TTL ±20%: разводит истечение ключей во времени.
Без джиттера 10k ключей истекают в одну секунду →
10k одновременных запросов к БД. С джиттером — размазано на 120 с.
"""
return int(base * (1 + random.uniform(-spread, spread)))
5.2 Событийная инвалидация
При изменении данных публикуется событие, потребители удаляют затронутые ключи. Точнее TTL, но требует надёжной доставки — то есть outbox и идемпотентности. Пропущенное событие означает вечно устаревший кэш, поэтому событийную инвалидацию всегда дополняют TTL как страховкой.
5.3 Версионирование ключей
Вместо удаления — смена ключа. Ключ строится как user:42:v7 или product:88:{updated_at}. При изменении версия растёт, старый ключ никто не читает и он вытесняется сам.
Изящный вариант — версия на группу: держим cache_epoch:catalog = 15, ключи строим как catalog:15:product:88. Инкремент эпохи мгновенно инвалидирует всю группу без перебора ключей. Это ровно то, что делает Rails с cache_key_with_version и generational caching (Rails Guides: Caching).
5.4 Порядок операций: удалять, а не обновлять
Классическая ошибка — при записи в БД обновлять кэш новым значением. Правильно — удалять запись из кэша (cache-aside invalidation). Причина в гонках: два конкурирующих обновления могут записать значения в кэш в порядке, обратном порядку записи в БД, и кэш навсегда останется рассогласованным. Удаление идемпотентно и не имеет такой проблемы.
Но и удаление не спасает от гонки «читатель против писателя»:
до истечения TTL. Инвалидация опоздала.
Окно узкое, но при высоком RPS оно реализуется ежедневно. Три способа лечения по возрастанию сложности:
- TTL как страховка — рассогласование живёт максимум
T. В 95% продуктов этого достаточно. - Отложенное второе удаление (delayed double delete): писатель удаляет ключ, потом через ~500 мс удаляет ещё раз, уже после того как все «в полёте» читатели успели записать старое.
- Лизы (leases) — механизм из memcached в Facebook: на промахе клиент получает токен, и записать значение может только владелец токена, а токен аннулируется любой инвалидацией. Описано в «Scaling Memcache at Facebook», NSDI 2013 — лучшая из существующих статей о кэшировании в проде, читать целиком.
6. Патологии кэша и их лечение
6.1 Cache stampede (thundering herd)
Горячий ключ истекает. Тысяча одновременных запросов промахивается и все идут в БД за одним и тем же значением. База складывается, ключ так и не записывается, эффект самоусиливается.
Лечение первого уровня — однопоточная перезагрузка: только один запрос идёт в БД, остальные ждут или получают старое значение.
import threading, time
from typing import Any, Callable
class SingleFlight:
"""Схлопывает конкурентные промахи по одному ключу в один вызов loader.
Аналоги: golang.org/x/sync/singleflight, Caffeine LoadingCache,
RedisTemplate + распределённый lock для межпроцессного случая.
"""
def __init__(self) -> None:
self._lock = threading.Lock()
self._inflight: dict[str, threading.Event] = {}
self._results: dict[str, Any] = {}
def do(self, key: str, loader: Callable[[], Any]) -> Any:
with self._lock:
event = self._inflight.get(key)
if event is None:
# Мы — лидер: выполняем загрузку сами
event = threading.Event()
self._inflight[key] = event
leader = True
else:
leader = False
if not leader:
event.wait(timeout=5.0) # ждём чужого результата
return self._results.get(key)
try:
value = loader()
self._results[key] = value
return value
finally:
with self._lock:
del self._inflight[key]
event.set()
self._results.pop(key, None)
Внутри одного процесса этого достаточно. Между процессами нужен распределённый замок: SET lock:user:42 <uuid> NX PX 3000 в Redis, и обязательно освобождение через Lua-скрипт со сверкой значения — иначе вы снимете чужой замок после своего таймаута.
Лечение второго уровня и более элегантное — вероятностное раннее истечение (XFetch). Идея: чем ближе к истечению и чем дороже пересчёт, тем выше шанс, что конкретный читатель обновит значение заранее. Из статьи Vattani, Chierichetti, Lowenstein, «Optimal Probabilistic Cache Stampede Prevention», VLDB 2015:
import math, random, time
def should_recompute(delta: float, expiry: float, beta: float = 1.0) -> bool:
"""delta — сколько секунд занял последний пересчёт значения,
expiry — абсолютное время истечения (monotonic),
beta > 1 делает обновление более агрессивным.
Возвращает True с вероятностью, растущей экспоненциально
по мере приближения к expiry. Дорогие ключи обновляются раньше дешёвых.
"""
now = time.monotonic()
return now - delta * beta * math.log(random.random()) >= expiry
Красота решения в том, что оно не требует ни блокировок, ни координации: каждый читатель принимает решение локально, а вероятность того, что обновят двое, ничтожна. Сложность — O(1), память — O(1) на ключ (нужно хранить delta).
6.2 Cache penetration — запросы за несуществующим
Атакующий (или сломанный клиент) запрашивает user:-1, user:-2… Такого пользователя нет, кэш не заполняется — каждый запрос уходит в БД.
Два лечения:
- Негативное кэширование: пишем в кэш маркер «не существует» с коротким TTL (30–60 с). Просто и работает.
- Фильтр Блума перед кэшем: структура на
mбит сkхеш-функциями даёт «точно нет» или «возможно есть» заO(k)и памятью ~10 бит на элемент при 1% ложных срабатываний. Ложноположительный ответ безвреден — просто сходим в БД зря.
6.3 Горячий ключ
Один ключ (курс валюты, товар с распродажи, ID знаменитости) получает столько запросов, что его шард в Redis упирается в CPU одного ядра — Redis однопоточен по обработке команд.
Приёмы, в порядке предпочтения:
- Локальный L1-кэш с TTL в 1–5 секунд. Горячий ключ по определению читается часто, значит локальный кэш почти всегда попадает. Самое дешёвое решение, срабатывает в большинстве случаев.
- Размножение ключа: храним
hot:rate:usd#0…hot:rate:usd#9, клиент читает случайную копию. Нагрузка размазывается по 10 слотам и, при консистентном хешировании, по разным узлам. - Реплики чтения кэш-кластера с чтением с ближайшей реплики.
6.4 Cache avalanche
Массовое одновременное истечение или падение всего кэш-кластера. БД получает 100% трафика, который никогда не рассчитывала выдержать — и падает, а вслед за ней падает всё. Это самый частый сценарий каскадного отказа с участием кэша.
Защиты: джиттер TTL (см. выше), ограничение конкурентности к БД (bulkhead), circuit breaker перед источником, и обязательно — нагрузочный тест «кэш выключен». Если ваша БД не переживает работу без кэша хотя бы на 20% трафика, у вас не оптимизация, а скрытая зависимость. Все эти механизмы подробно разобраны в следующей статье, «Устойчивость: circuit breaker, retry, bulkhead».
7. Вытеснение: что выбросить, когда память кончилась
Кэш конечен. Политика вытеснения определяет hit ratio при заданном объёме памяти — и разница между политиками достигает десятков процентов.
| Политика | Идея | Слабость | Стоимость |
|---|---|---|---|
| FIFO | выбрасываем самое старое | игнорирует популярность | O(1), минимум метаданных |
| LRU | выбрасываем давно не использованное | скан большого объёма вымывает весь кэш | O(1), список + хеш |
| LFU | выбрасываем редко используемое | «залипает» на старой популярности | O(1) с ведёрками |
| ARC | адаптивно балансирует recency и frequency | был под патентом IBM | O(1), 2× метаданных |
| W-TinyLFU | частотный эскиз + LRU-окно | сложнее в реализации | O(1), ~4 бита на ключ |
| S3-FIFO | три FIFO-очереди, отсев «одноразовых» | новизна, мало реализаций | O(1), дёшево |
Практический ориентир: W-TinyLFU (реализован в Caffeine) на реальных трассах даёт hit ratio, близкий к оптимальному Belady, и заметно выше LRU при том же объёме. Секрет — компактный частотный эскиз (Count-Min Sketch с обнулением-«старением»), который позволяет спросить «а этот кандидат вообще популярнее того, кого мы собираемся выбросить?» до допуска в кэш. Свежая альтернатива — S3-FIFO (SOSP 2023), которая достигает похожего качества тремя обычными FIFO-очередями и потому лучше ложится на многоядерность.
В Redis политика задаётся maxmemory-policy. Разумные умолчания: allkeys-lru для чистого кэша, volatile-ttl если в том же инстансе лежат данные без TTL, которые терять нельзя. noeviction (умолчание!) для кэша — почти всегда ошибка: при заполнении памяти записи начнут падать с ошибкой вместо вытеснения.
8. HTTP-кэширование: бесплатный уровень, который забывают
Перед тем как строить Redis, посмотрите, сколько трафика можно снять заголовками. Это буквально бесплатно.
Cache-Control: public, max-age=60,
stale-while-revalidate=600
ETag: "v42" E-->>B: 200 OK (+ те же заголовки) Note over B,E: 0–60 с: свежий ответ B->>B: отдаёт из локального кэша, сеть не задействована Note over B,E: 60–660 с: окно stale-while-revalidate B->>E: GET /api/catalog E-->>B: 200 OK (устаревший на 3 минуты) — мгновенно E->>O: фоновая ревалидация
If-None-Match: "v42" O-->>E: 304 Not Modified (тело не передаётся) E->>E: продлевает свежесть Note over B,O: Инвалидация по событию O->>E: PURGE по surrogate-key "catalog" E->>E: помечает все связанные объекты устаревшими
Что здесь важно понять:
ETag+304экономят трафик, но не round-trip. Пользователь всё равно ждёт полный RTT. Поэтомуmax-ageценнееETag.stale-while-revalidate(RFC 5861) — самый недоиспользуемый заголовок в вебе. Он превращает «медленный запрос раз в минуту» в «всегда быстро, обновление в фоне». Тот же принцип, что у refresh-ahead, но реализованный инфраструктурой.stale-if-error— отдавать устаревшее, если origin лежит. Бесплатная деградация.- Surrogate keys (теги кэша) — механизм Fastly/Cloudflare: помечаете ответ тегами (
Surrogate-Key: catalog product-88), а потом инвалидируете все ответы с тегом одной командой. Это событийная инвалидация на уровне CDN. Vary— тонкое место.Vary: Accept-Encodingнормально,Vary: User-Agentфрагментирует кэш до бесполезности.
Каноническое объяснение семантики — MDN: HTTP caching и RFC 9111. Про то, как это укладывается в дизайн API, — статья «Стили API».
9. Масштабирование данных, часть 1: репликация
Кэш решает проблему повторяющихся чтений. Репликация решает проблему объёма чтений вообще.
Схема стандартная: один лидер принимает записи, N реплик применяют его журнал и обслуживают чтения. Это ось X куба масштабирования, применённая к БД. Дёшево, поддерживается всеми СУБД из коробки — и приносит ровно одну, но принципиальную проблему: лаг репликации.
Нарушена гарантия read-your-writes
Лечения, по возрастанию цены:
- Sticky reads после записи. После записи пользователя в течение
2 × p99(лаг)направлять его чтения на лидера. Хранить факт в сессии или в куке. Дёшево и покрывает почти все жалобы. - Чтение по LSN. Приложение запоминает позицию журнала после записи (
pg_current_wal_lsn()в PostgreSQL) и требует от реплики догнать её. Строго, но требует поддержки в слое доступа к данным. - Монотонные чтения. Прикреплять пользователя к одной реплике, чтобы он хотя бы не видел, как время идёт назад.
Кэш и реплики взаимодействуют коварно: если наполнять кэш из реплики сразу после записи, вы закэшируете устаревшее значение и лаг репликации превратится из 80 мс в TTL кэша. Правило: прогрев кэша только с лидера, либо после подтверждённой репликации.
10. Масштабирование данных, часть 2: шардирование
Когда данные не помещаются на одну машину (по объёму, по IOPS или по скорости записи), остаётся ось Z: разделить данные между узлами. Реплики здесь не помогают — они все содержат всё.
Выбор ключа шардирования — главное решение
Ключ шардирования определяет всё дальнейшее: какие запросы будут дешёвыми, какие — катастрофическими, и можно ли будет вообще что-то поменять потом. Ошибка в ключе не исправляется — только полной перезаливкой.
Разбор квадрантов:
tenant_idдля B2B SaaS — почти всегда правильный ответ. Все запросы одного клиента живут в одном шарде, JOIN’ы работают. Риск: клиент-гигант перевешивает шард — лечится выделением ему персонального шарда.created_at— классическая ловушка. Все записи сегодняшнего дня идут в один шард: 100% нагрузки на запись на одну машину, остальные простаивают. Годится только для архивных данных с равномерным чтением по истории (тайм-серии, где шард = месяц и старые месяцы редко читают).hash(order_id)даёт идеальную равномерность и убивает любой запрос вида «все заказы клиента» — он превращается в fan-out по всем шардам с известной проблемой хвостовой латентности.
Три способа отобразить ключ на шард
Диапазонное разбиение. Шард владеет интервалом ключей (A–F, G–M, …). Плюс: диапазонные запросы (WHERE date BETWEEN) читают один шард. Минус: перекос почти гарантирован, а ребалансировка требует расщепления диапазонов. Так работают HBase, Bigtable, CockroachDB.
Хеш-разбиение. shard = hash(key) mod N. Плюс: идеальная равномерность. Минус: диапазонные запросы невозможны, а изменение N переселяет почти всё.
Каталог (directory-based). Отдельный сервис хранит карту «ключ → шард». Максимальная гибкость (можно переселить одного клиента вручную), цена — ещё один компонент на критическом пути, который надо кэшировать и который становится единой точкой отказа.
Консистентное хеширование
Наивное mod N при добавлении узла перемещает долю N/(N+1) ключей — при переходе с 3 узлов на 4 переезжает 75% данных. Для кэша это означает мгновенное падение hit ratio почти до нуля и, следовательно, лавину на БД. Консистентное хеширование (Karger et al., 1997) решает задачу: перемещается только 1/(N+1).
import bisect, hashlib
from typing import Iterable, Optional
class ConsistentHashRing:
"""Кольцо консистентного хеширования с виртуальными узлами.
Сложность:
lookup — O(log(N·V)) бинарным поиском по отсортированному кольцу
add/remove узла — O(V log(N·V))
память — O(N·V), при V=150 и N=100 это 15 000 int — десятки КБ
Дисбаланс нагрузки при V виртуальных узлах на физический ≈ O(1/√V):
V=1 даёт разброс в разы, V=150 — единицы процентов.
"""
def __init__(self, nodes: Iterable[str] = (), vnodes: int = 150):
self._vnodes = vnodes
self._ring: dict[int, str] = {}
self._sorted: list[int] = []
for node in nodes:
self.add(node)
@staticmethod
def _hash(key: str) -> int:
# blake2b быстрее md5 и даёт хорошее распределение;
# криптостойкость здесь не нужна, нужна равномерность
return int.from_bytes(hashlib.blake2b(key.encode(), digest_size=8).digest(), "big")
def add(self, node: str, weight: int = 1) -> None:
"""weight позволяет дать мощной машине больше токенов на кольце."""
for i in range(self._vnodes * weight):
point = self._hash(f"{node}#{i}")
self._ring[point] = node
bisect.insort(self._sorted, point)
def remove(self, node: str, weight: int = 1) -> None:
for i in range(self._vnodes * weight):
point = self._hash(f"{node}#{i}")
if self._ring.pop(point, None) is not None:
idx = bisect.bisect_left(self._sorted, point)
self._sorted.pop(idx)
def get(self, key: str) -> Optional[str]:
"""Первый виртуальный узел по часовой стрелке от hash(key)."""
if not self._sorted:
return None
h = self._hash(key)
idx = bisect.bisect_right(self._sorted, h) % len(self._sorted)
return self._ring[self._sorted[idx]]
def get_replicas(self, key: str, n: int) -> list[str]:
"""N различных физических узлов подряд — база для репликации,
как в Dynamo: реплики кладутся на следующие N узлов кольца."""
if not self._sorted:
return []
h = self._hash(key)
idx = bisect.bisect_right(self._sorted, h) % len(self._sorted)
result: list[str] = []
for offset in range(len(self._sorted)):
node = self._ring[self._sorted[(idx + offset) % len(self._sorted)]]
if node not in result:
result.append(node)
if len(result) == n:
break
return result
Проверка равномерности — обязательный тест, а не факультатив:
from collections import Counter
ring = ConsistentHashRing(["node-a", "node-b", "node-c", "node-d"], vnodes=150)
dist = Counter(ring.get(f"user:{i}") for i in range(200_000))
# node-a 50_412, node-b 49_780, node-c 50_101, node-d 49_707 — разброс < 1%
# Убираем узел: проверяем, что переехало ~1/4 ключей, а не 3/4
before = {f"user:{i}": ring.get(f"user:{i}") for i in range(200_000)}
ring.remove("node-c")
moved = sum(1 for k, v in before.items() if ring.get(k) != v)
print(moved / len(before)) # ≈ 0.25 — ровно доля ушедшего узла
Альтернативы, которые стоит знать: rendezvous hashing (HRW) — проще в реализации (O(N) на поиск, но без структуры кольца) и естественно поддерживает веса; jump consistent hash от Google — семь строк кода, O(log N), идеальная равномерность, но не умеет удалять произвольный узел из середины.
Фиксированное число партиций — приём, который решает ребалансировку
Самый практичный трюк, описанный у Клеппмана в «Designing Data-Intensive Applications» (глава 6): создать сильно больше партиций, чем узлов — скажем, 1024 партиции на 8 узлов — и раз и навсегда зафиксировать отображение «ключ → партиция». Ребалансировка тогда не трогает хеш-функцию вообще: она просто передвигает целые партиции между узлами.
Так устроены Elasticsearch (число primary shards фиксируется при создании индекса — и это самое частое место, где ошибаются), Kafka (число партиций топика), Riak, Couchbase. Цена решения — вы обязаны угадать порядок величины заранее: слишком мало партиций упрёт вас в потолок числа узлов, слишком много — накладные расходы на метаданные и файловые дескрипторы. Разумное правило: партиций в 10–100 раз больше, чем узлов на горизонте трёх лет.
Что ломается после шардирования
Честный список, который стоит прочитать до, а не после:
- JOIN между шардами не существует. Либо денормализация, либо два запроса и соединение в приложении, либо копия справочника на каждом шарде.
- Транзакции между шардами требуют 2PC (медленно, блокирующе) или саг — см. «Saga, распределённые транзакции, outbox и идемпотентность».
- Глобальная уникальность (
UNIQUE(email)) не обеспечивается локальным индексом. Нужен отдельный сервис-регистр уникальности или ключ шардирования, совпадающий с уникальным полем. - Автоинкрементные ID ломаются. Нужны Snowflake-подобные ID (время + номер узла + счётчик) или UUIDv7 — они, кстати, ещё и сортируются по времени, что важно для локальности в B-дереве.
- Вторичные индексы: локальный (по документу) требует fan-out на чтении; глобальный (по термину) требует распределённой записи. Оба варианта хуже, чем было.
COUNT(*)и агрегации становятся распределёнными задачами. Обычно ответ — заранее посчитанные счётчики, обновляемые событиями.- Ребалансировка — отдельный проект: двойная запись, копирование, сверка, переключение чтений, зачистка. Планируйте недели, не дни.
Именно поэтому шардирование — последнее средство. Сначала: индексы, кэш, реплики чтения, архивация холодных данных, вертикальное масштабирование, вынос BLOB’ов в объектное хранилище. Всё это дешевле на порядок.
11. Типичные ошибки
Кэшируют без измерений. Кэш ставят «потому что медленно», не выяснив, где именно медленно. Профилируйте: часто узкое место — N+1 запрос или отсутствующий индекс, и тогда кэш просто маскирует проблему, которая рванёт при следующем скачке нагрузки.
Не мониторят hit ratio. Кэш с hit ratio 20% — это чистый вред: вы платите сеть и память, а нагрузку не снимаете. Метрики-минимум: hit ratio по группам ключей, распределение TTL, объём вытеснений в секунду, p99 латентности кэша, доля stale-ответов.
Считают кэш опциональным, но проектируют так, что он обязателен. «Если Redis упадёт, просто пойдём в БД» — до первого падения Redis. Проверять надо учением: отключите кэш на 5% трафика в проде и посмотрите на графики БД.
Кэшируют персональные данные с общим ключом. Классическая уязвимость: GET /profile кэшируется на CDN без Cache-Control: private, и один пользователь видит профиль другого. Правило: любой ответ, зависящий от заголовка Authorization или куки — private, no-store по умолчанию.
Обновляют кэш вместо удаления. Разобрано в §5.4 — гонка приводит к вечно неверному значению.
TTL «побольше, чтобы наверняка». TTL — не параметр производительности, а контракт со стороны бизнеса: сколько секунд пользователь может видеть устаревшую цену. Спрашивайте у продукта, а не выбирайте сами.
Кладут в кэш большие объекты. Значение в 5 МБ забивает сеть, вызывает фрагментацию памяти и блокирует однопоточный Redis на десятки миллисекунд, ломая латентность всем остальным. Ограничение здравого смысла — сотни килобайт.
Шардируют по полю, которое меняется. Если ключ шардирования может измениться (например, region пользователя), то каждое такое изменение — это удалить-и-переселить между шардами, без транзакции. Ключ шардирования обязан быть неизменяемым.
Забывают про холодный старт. Деплой перезапускает 40 подов, у всех пустые L1-кэши, все одновременно идут в L2 и БД. Лечится постепенным rollout’ом, прогревом при старте и readiness-пробой, которая не пускает трафик до прогрева.
Отсутствие джиттера. Синхронное истечение, синхронные retry, синхронные крон-задачи в :00 — три источника рукотворных пиков.
12. Как это выглядит в проде
Facebook / Meta. Многоуровневый memcached на тысячи узлов; лизы против stampede и против устаревших записей; «холодный кластер» с прогревом из «тёплого» при вводе новых мощностей; региональная инвалидация через демона, читающего журнал MySQL. Всё описано в NSDI 2013 — это до сих пор обязательное чтение.
Netflix. EVCache — memcached с репликацией по зонам доступности: запись идёт во все AZ, чтение — из локальной, чтобы не платить за межзональный трафик и латентность. Плюс агрессивная деградация: если персонализация недоступна, показать всем одинаковый популярный ряд, а не ошибку.
Twitter/X. Знаменитый пример fan-out on write: лента формируется заранее и лежит в Redis, кроме пользователей с десятками миллионов подписчиков — для них применяется fan-out on read и слияние на чтении. Гибрид, потому что чистая стратегия ломается на хвостах распределения.
Discord. Переезд с Cassandra на ScyllaDB с шардированием сообщений по (channel_id, bucket), где bucket — временное окно. Плюс «слияние запросов» (request coalescing) в промежуточном сервисе на Rust: тысяча одновременных читателей одного канала превращается в один запрос к базе — тот же single-flight на уровне архитектуры. Технический блог Discord.
Stack Overflow годами работал на двух серверах SQL Server и вертикальном масштабировании с агрессивным многоуровневым кэшем, обслуживая сотни миллионов просмотров в месяц. Хорошее напоминание, что «как в Google» — не обязательное условие масштаба.
13. Практический порядок действий
Когда «система не тянет», работайте по этому списку сверху вниз и не перескакивайте:
- Измерьте. Где именно время: БД, сеть, CPU приложения, GC? Без профиля дальше не идти.
- Уберите глупость. N+1, отсутствующие индексы,
SELECT *по 200 колонок, сериализация JSON на каждый чих. - Включите HTTP-кэш.
max-age,ETag,stale-while-revalidate, CDN для статики и публичных ответов. Бесплатно. - Добавьте кэш приложения. Cache-aside + TTL с джиттером + single-flight. Измерьте hit ratio.
- Клонируйте stateless-инстансы (ось X). Дешевле любого другого шага.
- Реплики чтения. Плюс явная политика read-your-writes.
- Вертикально увеличьте БД. Часто это годы времени за небольшие деньги.
- Разделите по функциям (ось Y). Вынесите самый нагруженный домен.
- И только теперь — шардирование (ось Z). С заранее заготовленной картой партиций, планом ребалансировки и списком того, что сломается.
Между шагами 4 и 9 обычно лежат годы. Если вы прыгаете с 1 на 9, вы решаете не ту задачу.
Мини-итог
- Масштабируемость упирается не в число узлов, а в то, что у узлов общее: USL с ненулевым
βдаёт падение производительности при росте, а не плато. - Кэш — это осознанный обмен свежести на латентность. Величина обмена — TTL — это бизнес-решение, а не техническое.
- Уровни кэша образуют лестницу от процесса до браузера; каждый следующий быстрее на порядок и устареет сильнее. Считайте эффективную латентность по формуле
L = h·L_hit + (1−h)·(L_hit + L_miss)и помните, что важны промахи, а не попадания. - Инвалидация: удалять, а не обновлять; событийная инвалидация всегда подстрахована TTL; версионирование ключей избавляет от инвалидации вовсе.
- Патологии — stampede, penetration, горячий ключ, avalanche — имеют стандартные лечения (single-flight, XFetch, негативное кэширование, джиттер, L1 на горячих ключах). Они должны быть в коде до инцидента, а не после.
- Порядок масштабирования: X (клоны) → Y (функции) → Z (данные). Шардирование — последнее средство, потому что ключ шардирования выбирают один раз и навсегда.
- Консистентное хеширование перемещает
1/(N+1)ключей вместоN/(N+1); виртуальные узлы выравнивают нагрузку; фиксированное большое число партиций делает ребалансировку тривиальной.
Источники
- Мартин Клеппман, «Designing Data-Intensive Applications» — главы 5 (репликация) и 6 (партиционирование). Единственная книга, которую надо прочитать целиком по этой теме.
- «Scaling Memcache at Facebook», NSDI 2013
- Karger et al., «Consistent Hashing and Random Trees», STOC 1997
- Vattani et al., «Optimal Probabilistic Cache Stampede Prevention», VLDB 2015
- DeCandia et al., «Dynamo: Amazon’s Highly Available Key-value Store», SOSP 2007
- Neil Gunther, Universal Scalability Law
- Caffeine: Efficiency (W-TinyLFU) и S3-FIFO, SOSP 2023
- RFC 9111: HTTP Caching, RFC 5861: stale-while-revalidate
- Redis: key eviction policies
- Брендан Грегг, «Systems Performance», 2-е изд. — методология USE и закон Литтла на практике.
Что дальше
Кэш и шардирование делают систему быстрой и большой, но не делают её надёжной: чем больше узлов и зависимостей, тем выше вероятность, что что-то из них прямо сейчас лежит. Следующая статья — про то, как система переживает отказы своих частей: «Устойчивость: circuit breaker, retry, bulkhead, backpressure, деградация».