Архитектурные паттерны Кэширование и масштабирование: уровни, инвалидация, шардирование
0%

Кэширование и масштабирование: уровни, инвалидация, шардирование

Кэширование и масштабирование: уровни, инвалидация, шардирование

Есть два способа сделать систему быстрее: делать меньше работы или делать работу на большем количестве машин. Кэширование — это первое. Масштабирование — второе. Их объединяют в одну статью не по лени, а потому что это две стороны одного разговора: почти всегда правильный порядок действий — «сначала перестань считать одно и то же по десять раз, потом добавляй железо». Инженер, который начинает с шардирования базы, обычно через полгода обнаруживает, что 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-запроса через все уровни:

Разберём ступени содержательно.

Кэш браузера и 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). Причина в гонках: два конкурирующих обновления могут записать значения в кэш в порядке, обратном порядку записи в БД, и кэш навсегда останется рассогласованным. Удаление идемпотентно и не имеет такой проблемы.

Но и удаление не спасает от гонки «читатель против писателя»:

Окно узкое, но при высоком RPS оно реализуется ежедневно. Три способа лечения по возрастанию сложности:

  1. TTL как страховка — рассогласование живёт максимум T. В 95% продуктов этого достаточно.
  2. Отложенное второе удаление (delayed double delete): писатель удаляет ключ, потом через ~500 мс удаляет ещё раз, уже после того как все «в полёте» читатели успели записать старое.
  3. Лизы (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 однопоточен по обработке команд.

Приёмы, в порядке предпочтения:

  1. Локальный L1-кэш с TTL в 1–5 секунд. Горячий ключ по определению читается часто, значит локальный кэш почти всегда попадает. Самое дешёвое решение, срабатывает в большинстве случаев.
  2. Размножение ключа: храним hot:rate:usd#0hot:rate:usd#9, клиент читает случайную копию. Нагрузка размазывается по 10 слотам и, при консистентном хешировании, по разным узлам.
  3. Реплики чтения кэш-кластера с чтением с ближайшей реплики.

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, посмотрите, сколько трафика можно снять заголовками. Это буквально бесплатно.

Что здесь важно понять:

  • 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 куба масштабирования, применённая к БД. Дёшево, поддерживается всеми СУБД из коробки — и приносит ровно одну, но принципиальную проблему: лаг репликации.

Лечения, по возрастанию цены:

  1. Sticky reads после записи. После записи пользователя в течение 2 × p99(лаг) направлять его чтения на лидера. Хранить факт в сессии или в куке. Дёшево и покрывает почти все жалобы.
  2. Чтение по LSN. Приложение запоминает позицию журнала после записи (pg_current_wal_lsn() в PostgreSQL) и требует от реплики догнать её. Строго, но требует поддержки в слое доступа к данным.
  3. Монотонные чтения. Прикреплять пользователя к одной реплике, чтобы он хотя бы не видел, как время идёт назад.

Кэш и реплики взаимодействуют коварно: если наполнять кэш из реплики сразу после записи, вы закэшируете устаревшее значение и лаг репликации превратится из 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. Практический порядок действий

Когда «система не тянет», работайте по этому списку сверху вниз и не перескакивайте:

  1. Измерьте. Где именно время: БД, сеть, CPU приложения, GC? Без профиля дальше не идти.
  2. Уберите глупость. N+1, отсутствующие индексы, SELECT * по 200 колонок, сериализация JSON на каждый чих.
  3. Включите HTTP-кэш. max-age, ETag, stale-while-revalidate, CDN для статики и публичных ответов. Бесплатно.
  4. Добавьте кэш приложения. Cache-aside + TTL с джиттером + single-flight. Измерьте hit ratio.
  5. Клонируйте stateless-инстансы (ось X). Дешевле любого другого шага.
  6. Реплики чтения. Плюс явная политика read-your-writes.
  7. Вертикально увеличьте БД. Часто это годы времени за небольшие деньги.
  8. Разделите по функциям (ось Y). Вынесите самый нагруженный домен.
  9. И только теперь — шардирование (ось 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); виртуальные узлы выравнивают нагрузку; фиксированное большое число партиций делает ребалансировку тривиальной.

Источники


Что дальше

Кэш и шардирование делают систему быстрой и большой, но не делают её надёжной: чем больше узлов и зависимостей, тем выше вероятность, что что-то из них прямо сейчас лежит. Следующая статья — про то, как система переживает отказы своих частей: «Устойчивость: circuit breaker, retry, bulkhead, backpressure, деградация».

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов