Как принимать архитектурные решения в распределённых системах, которые работают с большими объёмами данных: что за чем следует, чем каждое решение платит и как это объяснить продукту. По мотивам «Designing Data-Intensive Applications».
Три ориентира, которыми стоит мерить любое архитектурное решение в высоконагруженной системе.
Надёжность, масштабируемость и удобство сопровождения тянут архитектуру в разные стороны: больше избыточности — дороже и сложнее; проще код — меньше гибкости при росте нагрузки. Хорошее архитектурное решение — не то, которое максимизирует все три цели сразу (это невозможно), а то, где вы явно называете, чем платите, и записываете это решение так, чтобы через год не гадать, почему сделали именно так.
Надёжность (reliability) — это способность системы продолжать корректно выполнять свою работу даже тогда, когда что-то внутри уже сломалось: диск умер, сеть моргнула, код упал с исключением. Ключевая мысль книги Клеппмана простая: отказ — это не аномалия и не «плохой день», а рабочее, ожидаемое состояние системы. Если вы проектируете архитектуру в предположении, что диски не ломаются и сеть не теряет пакеты, вы проектируете не для реального мира.
Внутри «надёжности» стоит различать три близкие, но разные вещи:
fault-free operation) — сколько система работает без единого отказа компонента. Измеряется через MTTF (mean time to failure, среднее время до отказа) — среднее время между двумя последовательными отказами одного и того же узла или диска.fault tolerance) — способность продолжать работать, когда отказ уже произошёл: одна реплика упала, но кластер отвечает как ни в чём не бывало. Здесь важен не MTTF самого узла, а MTTF системы в целом при том, что узлы внутри неё отказывают регулярно.repairability) — насколько быстро можно восстановить работоспособность после отказа. Мера — MTTR (mean time to repair, среднее время восстановления).Отсюда прямое определение доступности (availability): доля времени, когда система находится в рабочем состоянии, то есть MTTF / (MTTF + MTTR). Заметьте: доступность растёт двумя разными путями — реже ломаться или быстрее чиниться. На практике второй путь часто дешевле первого: чинить автоматизированно за 30 секунд проще, чем купить оборудование, которое не ломается вовсе.
redundancy): RAID, репликация, запасные узлы наготове.race condition), каскадная перегрузка одного сервиса из-за другого. Это самый коварный класс именно потому, что избыточность здесь не спасает: если баг в коде, то он одинаково сломает и основной узел, и все его копии одновременно — это систематическая, а не случайная ошибка.DROP TABLE, деплой не той версии. По данным разных отраслевых постмортемов люди — самая частая причина крупных, заметных пользователям сбоев, чаще, чем «железо» само по себе.Отдельно стоит упомянуть безопасность (security) как часть надёжности в широком смысле: защита от злонамеренных атак — это тоже про «продолжать работать правильно», просто источник неблагоприятных условий здесь не случайность, а противник, который специально ищет слабое место.
bulkhead, «переборка») — разбить систему на независимые ячейки (cell-based architecture) так, чтобы отказ одной ячейки не утопил остальные — как переборки на корабле не дают одному пробитому отсеку затопить весь корпус.backpressure (обратное давление) — когда потребитель не успевает, он явно сигналит источнику «помедленнее», а не тихо копит нарастающую очередь до взрыва по памяти.exponential backoff + jitter) — при ошибке клиент не долбит сервис немедленно, а ждёт всё дольше с каждой попыткой (1 с, 2 с, 4 с…). thundering herd, «эффект стада»). Джиттер размазывает повторные попытки во времени и убирает этот резонанс.chaos engineering) — намеренно вызывать отказы в контролируемых условиях (обычно не в проде «просто так», а через управляемые эксперименты), чтобы найти слабые места до того, как их найдёт реальный сбой.circuit breaker) — если зависимость стабильно отвечает ошибками, клиент временно перестаёт её вызывать вообще и сразу отдаёт быстрый отказ или запасной вариант, вместо того чтобы копить таймауты.checksums, end-to-end validation) — не доверять «канал передачи гарантированно доставит байты без искажений»; проверять контрольные суммы на границах, где данные реально важны.| Класс отказа | Типичный пример | Основная защита |
|---|---|---|
| Аппаратный | Диск, память, сетевой коммутатор | Избыточность, RAID, репликация, запасные узлы |
| Программный | Баг, утечка ресурсов, гонка, каскад перегрузки | Изоляция отказов, таймауты, предохранитель, канареечные релизы |
| Человеческий | Ошибка конфигурации, неверный деплой | Раздельные окружения, постепенный откат, автоматизация ручных шагов |
| Злонамеренный | Атака, попытка эксплуатации | Аутентификация, лимиты запросов, аудит, изоляция привилегий |
Доступность обычно измеряют долей времени простоя в год:
Важный нюанс, который легко упустить: SLA отдельного сервиса — это не то же самое, что опыт конечного пользователя. Если один пользовательский запрос по дороге проходит через 10 внутренних сервисов, каждый с доступностью 99.9%, итоговая вероятность, что все десять сработают, — примерно 0.99910 ≈ 99%, то есть в целом путь пользователя ощутимо менее надёжен, чем любой отдельный сервис в нём. Чем длиннее цепочка вызовов, тем дороже держать высокую доступность на каждом шаге.
Масштабируемость (scalability) — это не свойство «система быстрая» или «система медленная» само по себе. Это способность справляться с ростом нагрузки разумным приростом ресурсов. Чтобы говорить о масштабируемости осмысленно, нужно сначала точно назвать, что именно у вас растёт — это и называется параметром нагрузки.
load parameters)Нагрузка редко описывается одним числом «запросов в секунду». В разных системах узкое место создают разные величины:
| Тип системы | Параметр нагрузки | Пример |
|---|---|---|
| Веб-сервис / API | Запросов в секунду (RPS/QPS) | RPS на конкретный эндпоинт — общий RPS редко информативен |
| Лента / рассылка | Соотношение «чтение : запись» | В соцсетях чтения ленты обычно на порядки превышают число публикаций |
| Социальная сеть с подписками | Фан-аут (fan-out) — сколько записей или уведомлений порождает одно действие | Публикация одного твита должна быть видна у каждого подписчика |
| Key-value / кэш | Размер рабочего набора и плотность обращений к ключам | Соотношение попаданий и промахов кэша при заданном объёме RAM |
| Чат / совместное редактирование | Частота обновлений и целевая задержка доставки | Событий в секунду плюс требуемое время round-trip |
| Аналитика / batch | Пропускная способность конвейера | Строк в секунду через пакетную или потоковую обработку |
Классический пример из книги — почему у Twitter «болела» именно домашняя лента, а не публикация твитов сама по себе. У обычного пользователя порядка 350 подписок, но у отдельных знаменитостей — миллионы подписчиков. Здесь есть два принципиально разных подхода к тому, как собрать ленту:
Решение авторов Twitter — гибрид: для обычных пользователей твит толкается в кэши лент подписчиков заранее (дешёвое чтение), а твиты знаменитостей с огромным числом подписчиков в этот fan-out не идут — они подмешиваются в ленту в момент чтения, отдельным запросом. Общий урок шире Twitter: всегда спрашивайте не «сколько у нас RPS», а «что конкретно становится дороже при росте — размер списка подписчиков, размер fan-out, размер ответа — и как это ведёт себя при удвоении масштаба».
throughput) — сколько операций система обрабатывает за единицу времени. Естественная метрика для пакетных и потоковых конвейеров.latency) против времени отклика (response time) — время отклика — это всё время, которое видит клиент от отправки запроса до получения ответа, включая время обработки; задержка — это именно время ожидания в очереди, пока запрос ждёт своей обработки. На практике их часто путают, но разница важна при диагностике: если растёт время отклика, а не сама обработка, — вероятно, растёт очередь, а не сложность запроса.p50, p95, p99, p999) — p99 означает: 99% запросов отвечают быстрее этого значения, а 1% — медленнее. Среднее (avg) врёт почти всегда: распределение времени ответа обычно сильно смещённое — большинство запросов быстрые, но редкие запросы (упёрлись в GC-паузу, попали на перегруженный узел, ждали блокировку) отвечают в разы дольше, и именно эти редкие случаи вытягивают среднее вверх, маскируя при этом, что типичный опыт пользователя на самом деле быстрый. SLO стоит формулировать в перцентилях, например «p99 < 300 мс», а не «средняя задержка 80 мс» — второе ничего не говорит про то, каково 1% самых неудачливых пользователей.p99 с каждого из 50 серверов и посчитать их среднее — это математически неверно и даёт обманчиво оптимистичную картину. Правильный путь — собирать гистограммы распределения задержек с каждого узла и объединять именно гистограммы (например, через HDR-гистограммы, High Dynamic Range Histograms — компактную структуру, которая с высокой точностью хранит распределение по широкому диапазону значений), а перцентиль считать уже по объединённому распределению.p100) — это, по сути, время ответа самого медленного запроса за период, то есть «последний ответ»: метрика, крайне чувствительная к единичным выбросам, полезная для поиска аномалий, но плохая как SLO-цель.Формула из теории массового обслуживания: L = λ · W, где L — среднее число запросов, находящихся в системе одновременно, λ (лямбда) — скорость поступления запросов (например, запросов в секунду), а W — среднее время, которое запрос проводит в системе. Например, при 1000 запросов в секунду и среднем времени обработки 50 мс в системе одновременно находится в среднем около 50 запросов — именно столько параллельных «слотов» (потоков, соединений, воркеров) нужно держать наготове.
Практический вывод жёстче, чем кажется: если входящий поток λ устойчиво превышает скорость, с которой система успевает обслуживать запросы, очередь растёт без ограничения — и никакой «буфер побольше» это не лечит, он лишь откладывает момент, когда задержка станет неприемлемой. Единственные настоящие решения — увеличить скорость обслуживания (больше воркеров, быстрее код) или явно отбрасывать часть входящего потока (backpressure, троттлинг).
Есть два прямо противоположных способа дать системе больше мощности:
scale-up, вертикальное) — взять более мощную машину: больше ядер, больше RAM, быстрее диск. Плюс — простота: код не меняется, никакого распределённого состояния не появляется. Минус — есть физический потолок одного узла, и крупные машины стоят непропорционально дороже: машина с вдвое большим числом ядер обычно стоит больше, чем вдвое дороже, то есть экономическая эффективность падает по мере роста, хотя до определённого предела (одна машина вместо мелкого кластера) это часто всё равно проще и дешевле в эксплуатации.scale-out, горизонтальное) — добавить больше обычных машин и распределить нагрузку между ними. Плюс — нет теоретического потолка, можно расти постепенно. Минус — появляется распределённое состояние: нужно решать, как делить данные между узлами (шардирование, партиционирование — см. §3), как их синхронизировать, что делать при отказе одного из узлов. Это на порядок увеличивает инженерную сложность.На практике зрелые highload-системы обычно комбинируют: масштабируются вверх до разумного предела на одном узле (это дёшево в эксплуатации), а затем — вширь, когда предел достигнут.
Цель не в том, чтобы система держала абсолютно любой пик на полную мощность — это экономически бессмысленно (держать инфраструктуру под редкий пик 365 дней в году). Цель — предсказуемо ухудшать сервис, сохраняя критичные для бизнеса функции, вместо неконтролируемого падения всего сразу.
| Паттерн | Суть | Чем платим |
|---|---|---|
| Кэширование горячих данных | Часто читаемые данные держим в быстрой памяти, а не бьём по основной БД на каждый запрос | Риск отдать неактуальные данные; нужно продумать TTL и инвалидацию (см. §5) |
| Ограничение очереди и backpressure | Буфер входящих запросов ограничен по размеру; при переполнении — явно отклонять новые запросы, а не копить бесконечно | Нужно решить, где именно ставить границу и что вернуть клиенту при отказе |
| Выборочный отказ от необязательных функций | Временно отключить рекомендации, аналитику, «похожие товары», оставив только критичный путь (например, оформление заказа) | Часть функциональности видимо недоступна пользователю |
| Последовательная деградация | Ступенчато снижать качество ответа: сначала полный расчёт, потом приблизительный, потом статичный fallback | Нужны заранее подготовленные «упрощённые» пути ответа для каждого уровня |
| Feature flags | Возможность на лету выключить дорогую фичу без деплоя, прямо во время инцидента | Дополнительный слой конфигурации, который нужно поддерживать и тестировать |
| Троттлинг клиентов | Ограничить частоту запросов от одного клиента/ключа API, чтобы один «шумный» потребитель не съел ресурс для всех остальных | Легитимные клиенты с всплеском активности тоже могут попасть под лимит |
Если один запрос пользователя внутри дёргает N независимых зависимостей параллельно (типичный fan-out на микросервисы) и каждая зависимость сама по себе быстра в 99% случаев (p99), вероятность, что хотя бы одна из N зависимостей окажется той самой медленной 1%-й, растёт с числом зависимостей примерно как 1 − 0.99^N. При N = 100 это уже около 63% — то есть почти каждый запрос ловит чей-то «хвост», даже если каждая отдельная зависимость сама по себе надёжна. Именно это наблюдение объясняет, почему одиночная медленная зависимость превращается в системную проблему при высоком fan-out, и почему в таких системах применяют фоновые дублирующие запросы (hedged requests): отправить второй, резервный запрос к той же зависимости (или к другой её реплике) после небольшой задержки, если первый ответ не пришёл, и брать в работу тот, что ответит раньше.
Стоимость системы не равна стоимости её первоначальной разработки. Основную часть расходов за жизненный цикл системы обычно составляет всё, что происходит после первого релиза: поддержка, доработки, устранение багов, адаптация к новым требованиям. Клеппман сводит удобство сопровождения (maintainability) к трём принципам:
operability) — насколько легко команде эксплуатации держать систему в рабочем состоянии. На практике это означает: понятные метрики и логи, предсказуемое поведение при типичных операциях (деплой, откат, масштабирование), отсутствие «магии», которую понимает только один человек в команде.simplicity) — насколько легко новому человеку понять, что происходит в системе. Здесь полезно различать существенную сложность (она неизбежно вытекает из самой задачи — нельзя упростить распределённый консенсус до тривиальности) и случайную сложность (она добавлена самой реализацией — неудачной абстракцией, наслоением временных решений, случайными связями между модулями). Бороться нужно именно со второй: хорошая абстракция скрывает случайную сложность, оставляя на виду только существенную.evolvability) — насколько легко менять систему под новые требования. Слабая связанность модулей и явные, стабильные интерфейсы между ними облегчают изменения; скрытая, «невидимая» логика, разбросанная по хранимым процедурам и триггерам базы данных, обычно затрудняет — такую логику трудно тестировать, версионировать и увидеть при чтении кода приложения.Абстракции — главный инструмент простоты: API скрывает детали реализации сервиса, декларативный запрос (SQL) скрывает от разработчика, как именно движок базы данных будет искать данные, оставляя только описание что нужно найти — в отличие от императивного кода, который явно диктует порядок действий. Чем выше уровень абстракции, тем проще рассуждать о системе, но тем меньше прямого контроля над деталями исполнения.
Данные со временем эволюционируют: меняются схемы таблиц, появляются новые версии API. Здесь важны понятия обратной и прямой совместимости (backward/forward compatibility) — старый код должен уметь работать с новыми данными и наоборот, иначе любое изменение схемы превращается в синхронный релиз всех сервисов сразу. Подробный разбор форматов сериализации и правил совместимости — в §2.
На уровне организации удобство сопровождения — это не только код, но и процессы: понятный график дежурств (on-call), разбор инцидентов без поиска виноватого (blameless postmortem — цель разбора не найти, кто ошибся, а понять, что в системе позволило одной ошибке привести к сбою), автоматические тесты, которые ловят регрессии до продакшена, и наблюдаемость системы в целом (метрики, логи, трассировка — подробнее в §6).
error rate) — доля запросов, завершившихся с ошибкой.SLO (service level objective) — внутренняя цель, например «p99 < 300 мс, доступность 99.9% за 30 дней»; SLA (service level agreement) — то же самое, но зафиксированное как формальное обязательство перед клиентом, часто с денежными последствиями за нарушение. Error budget — допустимый объём «плохого» поведения за период, вытекающий из SLO. Если цель — 99.9% доступности за месяц, бюджет ошибок — те самые ~43 минуты простоя, которые можно «потратить». Пока бюджет не исчерпан, команда может рисковать (быстрее катить фичи, экспериментировать); когда бюджет на исходе — приоритет смещается на стабильность. Это удобный общий язык между продуктом и разработкой: не абстрактный спор «стабильность против скорости», а конкретное число, которое обе стороны понимают одинаково.MTTR (среднее время восстановления) — насколько быстро команда возвращает систему в рабочее состояние после обнаружения проблемы.change failure rate) — какой процент деплоев привёл к сбою или откату; хороший индикатор качества процесса тестирования и релизов, не только качества кода.Выбор модели данных и движка хранения определяет, какие запросы будут дешёвыми, а какие — практически невозможными. Переделать этот выбор позже почти всегда дороже, чем сделать его правильно сразу.
Модель данных — это не вопрос вкуса «SQL или NoSQL». Это заранее сделанная ставка на то, какие операции над данными будут частыми и дешёвыми, а какие — редкими или дорогими. Механизм хранения внутри движка (LSM-дерево или B-дерево) — вторая такая же ставка, только уже про то, во что вы платите на диске: скорость записи или предсказуемость чтения. Обе ставки трудно переиграть после того, как система уже накопила данные и трафик.
| Модель | Когда выбирать | Чем платите |
|---|---|---|
| Реляционная (SQL) | Данные с выраженными связями «многие-ко-многим», где нужны соединения (joins) и ссылочная целостность важнее гибкости схемы |
Схема должна быть определена заранее; изменение схемы на большой таблице под нагрузкой требует осторожных миграций |
| Документная (NoSQL) | Данные естественно ложатся в один документ (профиль пользователя, заказ со всеми позициями), читаются и пишутся целиком, без частых соединений между коллекциями | Слабая или отсутствующая поддержка соединений; связи «многие-ко-многим» приходится решать денормализацией и дублированием данных, а не ссылками |
| Графовая | Главный паттерн доступа — обход связей произвольной глубины: социальные связи, рекомендации, маршруты, модели прав доступа | Горизонтальное масштабирование трудное: разрезать граф на части (партиционировать) так, чтобы связанные узлы оставались рядом, — отдельная нерешённая в общем виде задача |
| Колоночная (OLAP) | Аналитические запросы — агрегации по огромному числу строк, но по небольшому числу столбцов сразу | Плохо подходит для точечного доступа «дай мне одну запись по id» и для частых обновлений отдельных строк; требует отдельного процесса загрузки данных из операционной БД |
| Key-value | Точечный доступ по известному ключу: кэш, счётчики, сессии, часто как быстрый слой поверх основной БД | Никаких сложных запросов, только доступ по ключу; нужно продумать инвалидацию и, если хранилище не durable, — устойчивость к потере данных при перезапуске |
schema-on-write): попытка вставить «мусорную» запись без обязательного поля просто провалится с ошибкой — база сама гарантирует целостность. Документная база в типичном режиме принимает почти всё (schema-on-read): записать можно что угодно, но тогда разбираться с отсутствующими или лишними полями обязано само приложение при чтении. Разница напрямую влияет на то, где будет болеть миграция данных: в реляционной модели вы платите болью один раз, в момент деплоя миграции; в документной — платите постоянно, кодом на защитную проверку полей на каждом чтении, зато без блокирующей миграции всей таблицы сразу.
SQL — декларативный язык (declarative): вы описываете что хотите получить («все заказы дороже 1000 рублей за последний месяц»), а как это достать — решает оптимизатор запросов внутри движка. Императивный подход — это когда код сам явно задаёт порядок обхода данных: «пройди по этому индексу, для каждой записи проверь условие, собери результат». Декларативность важна для производительности именно потому, что оставляет движку свободу маневра: если появился более удачный индекс или изменилась статистика данных, оптимизатор может переписать план выполнения запроса без изменения самого запроса — а с императивным кодом вы вручную переписываете логику обхода при каждом таком изменении.
Язык модели данных и язык запросов — не обязательно одно и то же: документные базы иногда предоставляют SQL-подобный интерфейс для запросов по вложенным полям; GraphQL может стоять поверх реляционной базы как единый декларативный слой запросов для клиентов, не показывая наружу, что под ним на самом деле обычные таблицы и joins.
Для highload это не деталь, а критичный вопрос: сериализация (serialization) — превращение структуры данных в поток байт для передачи по сети или хранения на диске — происходит на каждом межсервисном вызове и в каждой записи в лог событий. Форматы делятся на текстовые (JSON, XML — читаемые человеком, но многословные и без встроенной схемы) и бинарные (Protocol Buffers, Avro, Thrift, MessagePack — компактнее и быстрее парсятся, но требуют инструментов, чтобы прочитать «глазами»).
Главная практическая проблема — совместимость схем при эволюции данных: поле переименовали, поле удалили, добавили новое обязательное поле. Здесь работают два понятия:
backward compatibility) — новый код может прочитать данные, записанные старым кодом.forward compatibility) — старый код может прочитать данные, записанные новым кодом (это важно, когда откатываете деплой или когда разные сервисы обновляются не одновременно).Бинарные форматы со схемой (Avro, Protobuf, Thrift) решают это через нумерацию полей и правило «новое поле должно быть опциональным» — старый код просто игнорирует поле, номер которого не знает, а не падает с ошибкой. Отдельный практический приём — хранить идентификатор схемы (schema id) в начале самого сообщения, а полную схему держать в отдельном реестре схем (schema registry): тогда сообщение можно кодировать компактно, не таская описание схемы в каждом байтовом пакете.
| Формат | Схема | Практический вывод |
|---|---|---|
| JSON | Нет строгой схемы (JSON Schema — опционально и снаружи формата) | Читаемый человеком, удобен для отладки и публичных API, но занимает больше места и не проверяет типы сам по себе — проверка целиком на приложении |
| Protocol Buffers | Обязательная схема (.proto-файл), кодогенерация классов | Компактный бинарный формат, быстрое (де)сериализация; удобен, когда обе стороны компилируют один и тот же IDL-файл в свой код |
| Avro | Схема хранится отдельно от данных (часто в реестре схем) | Данные сами по себе крайне компактны, так как без схемы не самодостаточны; удобно для потоковых систем типа Kafka, где схему можно версионировать централизованно |
Ещё одно практическое правило highload-моделирования: между документами лучше использовать явные ссылки на внешние идентификаторы, а не вкладывать полную копию связанной сущности внутрь документа. Вложенная копия быстрее читается за один запрос, но при обновлении исходной сущности (например, у пользователя поменялось имя) нужно находить и обновлять все документы, куда эта копия была вложена — источник аномалий обновления (update anomalies), когда часть копий обновили, а часть — забыли.
Индекс (index) — это дополнительная структура данных, которая ускоряет поиск нужных записей за счёт того, что при каждой записи приходится обновлять не только сами данные, но и сам индекс. Это фундаментальный компромисс любого индекса: чем быстрее с ним читать, тем дороже с ним писать, и наоборот. Разные механизмы хранения — это разные точки на этой шкале.
Поведение движка хранения объясняется физикой дисков. У классического HDD случайное чтение (перемещение головки к произвольному месту на пластине) занимает порядка 10 мс, тогда как последовательное чтение идёт со скоростью порядка 1 Гбайт/с — разница на несколько порядков. SSD убирает механическую составляющую: случайное чтение занимает порядка 0.1 мс, то есть примерно в 1000 раз быстрее, чем у HDD, но последовательный доступ всё ещё выгоднее случайного даже на SSD. У SSD есть своя особенность: они любят пакетные последовательные записи большими блоками и деградируют от частой перезаписи одних и тех же малых участков — это называют усилением записи (write amplification): логически вы записали немного данных, а физически контроллер SSD вынужден перезаписать и переместить куда больше, из-за того, как организована флеш-память внутри.
B-дерево разбивает базу данных на страницы фиксированного размера — обычно от 4 до 16 КиБ, часто именно 16 КиБ — и организует их в древовидную структуру, где поиск по ключу — это спуск от корня к листу через несколько страниц. Обновление происходит на месте: страница читается целиком, изменяется в памяти и перезаписывается на диск целиком же. Для надёжности при этом используется журнал предзаписи (WAL, write-ahead log, иногда называют redo-логом): перед изменением самой страницы в журнал дописывается запись о том, что сейчас будет сделано, — если машина упадёт посреди записи страницы, при перезапуске можно доиграть журнал и восстановить страницу в консистентное состояние. Это же даёт point-in-time recovery — возможность восстановить состояние базы на конкретный момент в прошлом, проигрывая журнал.
LSM-дерево (log-structured merge-tree) устроено принципиально иначе. Новые записи сначала попадают в отсортированную структуру в памяти (memtable), одновременно дублируясь в журнал предзаписи на диске для надёжности на случай падения до того, как memtable сброшен. Когда memtable заполняется, он целиком сбрасывается на диск одним последовательным блоком — новым сегментом, отсортированным по ключу (SSTable, sorted string table). Со временем сегментов накапливается много, и фоновый процесс слияния (compaction) объединяет несколько сегментов в один, отбрасывая устаревшие и удалённые записи. Запись при этом почти всегда — последовательный append, что и делает LSM быстрым на запись.
Плата за это — три вида «усиления» (amplification):
bloom filter — компактная вероятностная структура, которая по одной быстрой проверке в памяти почти всегда точно говорит «этого ключа точно нет в данном сегменте», позволяя не читать сегмент с диска вообще, если ключа там гарантированно нет.tombstones), которые физически занимают место, пока их не вычистит следующий цикл слияния.| Критерий | LSM-деревья | B-деревья |
|---|---|---|
| Скорость записи | Быстрее — последовательный append, амортизированная стоимость compaction | Медленнее — каждая запись требует произвольного чтения-изменения-записи страницы |
| Скорость чтения | Обычно медленнее и менее предсказуема — нужно проверить несколько уровней (спасают bloom-фильтры) | Быстрее и стабильнее — один детерминированный путь по дереву |
| Усиление записи | Выше — данные переписываются несколько раз за жизненный цикл через compaction | Ниже, но каждая запись случайная, а не последовательная |
| Усиление чтения | Выше без bloom-фильтров, умеренное с ними | Минимальное — фиксированная глубина дерева |
| Усиление пространства | Выше — дубли версий и tombstones до compaction | Ниже, хотя фрагментация страниц после разбиений (page splits) тоже отнимает место |
| Сжатие на диске | Лучше — плотное последовательное хранение | Хуже — частично заполненные страницы после разбиений |
| Кто использует | LevelDB, RocksDB, HBase, Cassandra, движок индексации Lucene | InnoDB (MySQL), PostgreSQL, SQLite, Oracle |
p99.page split) при её переполнении оставляет частично заполненные страницы, что со временем приводит к фрагментации.Помимо выбора между LSM и B-деревом, есть более тонкие решения на уровне самого индекса:
(user_id, created_at)): составной ключ ускоряет запросы, которые фильтруют или сортируют именно по этой комбинации полей в этом порядке, но бесполезен для запросов по одному лишь второму полю без первого.covering index) — индекс, который включает все поля, нужные для ответа на запрос, так что движку не нужно ходить в саму таблицу за дополнительными данными — ценой того, что индекс становится больше и дороже в поддержке при каждой записи.Колоночное хранение (column-oriented storage) для аналитических нагрузок (OLAP) переворачивает физическое размещение данных: вместо того чтобы хранить подряд все поля одной строки, все значения одного столбца хранятся подряд. Это даёт два эффекта. Во-первых, аналитический запрос, который трогает только 3 столбца из 50, физически читает с диска только эти 3 столбца, а не всю строку целиком. Во-вторых, значения одного столбца обычно намного более однородны, чем значения строки (например, столбец «страна» содержит лишь несколько десятков уникальных значений на миллионы строк), поэтому такие столбцы сжимаются в разы эффективнее — часто через словарное кодирование (заменить повторяющиеся значения короткими числовыми кодами) в сочетании с битовыми картами (bitmap) для быстрой фильтрации по этим кодам.
Вывод простой и важный: OLTP (частые точечные транзакции по одной-двум строкам) и OLAP (агрегации по миллионам строк, но по узкому набору столбцов) — это принципиально разные паттерны доступа к диску, и оптимальные структуры хранения для них противоположны. Смешивать их в одной базе без разбора — почти гарантированный путь к тому, что аналитический запрос положит продакшн-нагрузку OLTP. Гибридные системы (HTAP, hybrid transactional/analytical processing) существуют, но применять их стоит осторожно и только когда реально нужна аналитика над самыми свежими данными без задержки репликации — в остальных случаях простое разделение на операционную БД и отдельное аналитическое хранилище надёжнее.
За операционной базой почти всегда стоят производные индексы — например, полнотекстовый поиск через Lucene или отдельные колоночные витрины для отчётности. Это те самые «производные данные», отделённые от системы записи (system of record), про которые подробно — в §4.
Современный движок хранения редко пишет напрямую на «голый» диск — он работает поверх файловой системы или, всё чаще, поверх объектного хранилища вроде S3. Это отдельный архитектурный слой, который стоит проектировать осознанно.
Идея разделения вычислений и хранения (separation of compute and storage) в том, что слой, который обрабатывает запросы (вычислительные узлы), и слой, который физически держит данные (объектное хранилище или распределённая файловая система), масштабируются независимо друг от друга. Лог-структурированное хранение (тот же принцип, что у LSM-деревьев — последовательные неизменяемые сегменты) хорошо ложится на объектные хранилища именно потому, что объектные хранилища дружелюбны к последовательной записи больших неизменяемых блоков и плохо подходят для частых изменений «на месте» — а лог-структурированный подход и не требует изменений на месте.
Что это даёт: хранение становится дешёвым и почти неограниченно масштабируемым по объёму, а новые вычислительные реплики можно поднимать быстро, потому что им не нужно копировать данные — они просто подключаются к тому же общему хранилищу. Чем платите: задержка обращения к объектному хранилищу выше, чем к локальному диску, консистентность между вычислительными узлами становится отдельной задачей (кто и когда видит только что записанные данные), а сами запросы к объектному хранилищу обычно тарифицируются по числу операций и объёму трафика, то есть паттерн доступа напрямую влияет на счёт от облачного провайдера.
Практический вывод: разделяйте хранение и вычисления, если вам нужно масштабировать объём данных и нагрузку на них независимо друг от друга — например, когда данные растут постоянно, а пики вычислительной нагрузки короткие и предсказуемые. Если же нагрузка стабильна и предсказуема, а задержка критична, локально прикреплённое хранилище с прямым доступом к диску часто остаётся проще и быстрее.
system of record, источник истины) от производных индексов и кэшей, которые всегда можно перестроить из источника (см. §4).
Репликация · партиционирование · транзакции · CAP/PACELC · консенсус.
Распределённость добавляет три фундаментальные проблемы: сетевую задержку нельзя отличить от отказа, часы на узлах не совпадают, а порядок событий не задан сам собой. Почти каждая гарантия согласованности — это цена за иллюзию «одной машины»: больше задержка, меньше доступность или выше сложность.
Репликация — это несколько копий одних данных. Она нужна для трёх задач: обслуживать больше чтений, переживать отказ узла и держать копию ближе к пользователю. Но копии надо синхронизировать.
Свет в оптоволокне идёт примерно 200 000 км/с. Москва—Франкфурт даёт RTT порядка 30–40 мс. Если запись подтверждают три узла на разных континентах, только сетевые поездки могут добавить к фиксации порядка 100 мс и больше. Поэтому «синхронно везде» почти никогда не означает «быстро».
| Схема | Суть | Плюсы | Минусы | Когда применять |
|---|---|---|---|---|
| Single-leader | Один лидер принимает записи и рассылает их репликам. | Простой порядок, понятные транзакции. | Лидер — узкое место; нужен failover. | Начальный и основной выбор для большинства систем. |
| Multi-leader | Несколько узлов принимают записи независимо. | Запись ближе к ДЦ и офлайн-клиентам. | Конфликты и сложное слияние. | Несколько ДЦ с записью или офлайн-редактирование. |
| Leaderless | Чтение и запись идут в кворумМинимальное большинство узлов, достаточное для решения; обычно более половины. | Высокая доступность записи. | Слабее гарантии; конфликты надо разрешать. | Dynamo/Cassandra-подобные сценарии, где доступность важнее свежести. |
Клиент отправляет запись лидеру. Лидер добавляет её в журнал упреждающей записи (WAL), затем рассылает изменение репликам. Новый лидер выбирается через консенсус (см. §3.5), а старого сначала отсекают fencing-механизмом: его токен эпохи становится недействительным, и хранилища отвергают его поздние записи. Без fencing «бывший лидер» после сетевого разрыва может продолжить писать и испортить данные.
Replica lag — отставание реплики от лидера: по времени, номеру WAL или позиции журнала. Мониторьте p95/p99 lag, возраст самой старой записи, число реплик вне SLA; задайте порог, после которого реплику исключают из чтений.
Если реплика отстаёт, пользователь на мгновение увидит старое состояние. Это не обязательно ошибка — это выбранная цена скорости.
| Гарантия | Человеческий симптом | Практическое решение |
|---|---|---|
| Read-your-writes | Поставили лайк под своим постом, обновили страницу — лайк исчез. | После записи читать с лидера или реплики не старше нужной версии. |
| Monotonic reads | Переключились между репликами и увидели цену ниже, чем секунду назад. | Закрепить пользователя за одной репликой или передавать version stamp. |
| Consistent prefix | Лента показывает ответ раньше сообщения, на которое он отвечает. | Читать причинно полный префикс; использовать causality token. |
| Write integrity | Сервис подтверждает действие, не видя обязательной предыдущей записи. | Проверять зависимость на лидере или передавать причинную метку. |
Простые решения из книги: направлять чтения своих записей на лидер/свежую реплику, закреплять пользователя за одной репликой, передавать causality token или version stamp — номер версии, после которой читать уже нельзя.
На часах узлов нельзя бездумно строить last write wins: NTP может сдвинуть часы, и более старая запись окажется «позже». LWW с временем узла проще, но тоже может потерять обновление. Векторные часы хранят для каждой реплики версии, которые она уже видела: по векторам можно отличить «новее» от «две записи возникли параллельно». Дальше конфликт решает бизнес-логика.
CRDT — структуры данных с заранее определённым слиянием: счётчики с операциями +/-, множества с правилом «победитель получает всё», регистры с несколькими победителями. Независимо от порядка доставки итог одинаков. Любую запись делайте идемпотентной: храните ID операции и отбрасывайте повтор.
При n=3, w=2, r=2 условие w+r>n выполнено: два узла записи и два чтения пересекаются. Но две записи могут прийти параллельно, read-repair может ещё не закончиться, а sloppy quorum отправит данные на временные узлы. Поэтому N/R/W в Dynamo или Cassandra — настройка вероятности свежего чтения, а не автоматически линейная согласованность.
Шард — небольшой кусок базы. Цель — чтобы обычная операция попадала в один шард: тогда добавление узлов увеличивает объём и пропускную способность без постоянных межшардовых запросов.
| Стратегия | Суть | Плюсы | Минусы | Когда |
|---|---|---|---|---|
| Диапазоны | Ключи 0–999, 1000–1999 и т. п. | Быстрые диапазонные запросы. | Соседние ключи могут перегреть шард. | Время, журналы, аналитика. |
| Хеш ключа | hash(key) выбирает шард. | Равномерная нагрузка. | Диапазонный поиск по исходному ключу дорог. | Равномерные точечные запросы. |
| Составной ключ | Первая часть — шард-ключ, вторая — ключ внутри шарда. | Баланс и диапазоны внутри владельца. | Нужно заранее знать шаблон запросов. | Например, пользователь → его события. |
Шард-ключ должен иметь высокую кардинальность и равномерно распределять запросы. «Знаменитость» в социальной сети может получать в тысячи раз больше обращений, чем обычный пользователь, и перегреть один шард даже при хорошем хеше.
00..99 и писать в 100 под-ключей; чтение собирает их.| Индекс | Как устроен | Запись | Чтение |
|---|---|---|---|
Локальный (term-partitioned) | Каждый шард индексирует только свои документы. | Быстрая, локальная. | Scatter/gather по всем шардам. |
Глобальный (document-partitioned) | Индекс распределён по значению поля и знает владельца документа. | Нужно обновить документ и индекс; возможна рассинхронизация. | Быстрый поиск в нужной части. |
hash mod N плох: при изменении числа шардов почти каждый ключ получает новый адрес. Лучше consistent hashing с виртуальными шардами или фиксированное число логических сегментов: узлы получают и передают сегменты, а не пересчитывают все ключи. Redis Cluster использует слоты; похожий подход применяют VK и Couchbase. HBase/Cassandra могут динамически делить большие диапазоны. Переезд запускается автоматически по метрикам либо оператором в окно обслуживания; в обоих случаях нужны лимиты скорости миграции и контроль репликации.
Ограничения: межшардовые транзакции дороже; счётчик на одном ключе создаёт hotspot; «все ключи пользователя в одном шарде» часто лучше для профиля и заказов, потому что сохраняет локальность операции.
Транзакции дают удобную абстракцию: несколько объектов меняются как одно целое, а читателю не нужно самому угадывать, что происходит при параллельной записи. Многие NoSQL отказались от широких транзакций ради скорости, простоты масштабирования и работы через разделы сети.
Двухфазная фиксация (2PC) сначала просит участников подготовиться, затем отправляет всем commit. Если координатор падает между фазами, участники остаются в in-doubt: они не знают, фиксировать или откатывать, и могут блокировать ресурсы. 2PC — дорогая комбинация согласования и готовности, а не магия, делающая сеть надёжной.
В MVCC каждая транзакция читает снимок данных; читатели не блокируют писателей. Но снимок сам по себе не запрещает все гонки.
| Уровень | Что даёт | От чего не защищает |
|---|---|---|
| Read committed | Не читает чужие незавершённые записи. | Повторное чтение может измениться; фантомы. |
| Snapshot isolation / MVCC | Стабильный снимок, читатели не блокируют писателей. | Lost update, write skew, некоторые фантомы. |
| Serializable | Результат как при последовательном выполнении. | Почти ни от чего, но платит производительностью и abort'ами. |
| Гонка | Бытовой пример | Лечение |
|---|---|---|
| Dirty read | Видим перевод, который потом откатился. | Read committed или выше. |
| Nonrepeatable read | Баланс изменился между двумя чтениями. | Снимок/MVCC. |
| Phantom | В проверяемом диапазоне внезапно появился новый заказ. | SSI, блокировка диапазона, материализованный конфликт. |
| Lost update | Два оператора правят остаток, последний затирает первого. | Атомарный upsert, SELECT FOR UPDATE, CAS или версия. |
| Write skew | Два дежурных видят второго коллегу и оба уходят: никого не осталось. | SSI, блокировка общей строки, уникальное ограничение, материализация конфликта. |
Serializable через 2PL блокирует данные и рискует deadlock; через SSI отслеживает зависимости и прерывает одну из конфликтующих транзакций (так работает PostgreSQL SERIALIZABLE). Очередь, которая последовательно обрабатывает операции одного ключа, материализует конфликт. Calvin и Silo используют логический порядок/версии, уменьшая число блокировок.
Межшардовые транзакции редки и дороги. Вместо них используйте сагу с компенсирующими действиями, идемпотентность и события; проектируйте так, чтобы одна операция жила в одном шарде.
Деньги, инвентарь и лимиты — serializable либо явная блокировка. Обычные ленты и профили — snapshot плюс версии/идемпотентные операции. Не растаскивайте одну транзакцию по шардам без крайней необходимости.
CAP говорит: при разделе сети (P) нельзя одновременно сохранить линейную согласованность (C) и полную доступность (A). CP-система в раздел откажет части запросов, AP-система ответит, возможно устаревшими данными. Формула «выбрать 2 из 3» вводит в заблуждение: CA без разделов сети — не полезная категория, потому что разделы случаются.
PACELC уточняет цену в обычной работе: если P — выбираем A или C; иначе E — выбираем меньшую задержку L или согласованность C.
| Система/настройка | Типично | Оговорка |
|---|---|---|
| DynamoDB, Cassandra | PA/EL | Настройки чтения и кворума меняют поведение. |
| MongoDB, HBase | PA/EC или PC/EC | Зависит от режима репликации и чтения. |
| ZooKeeper, etcd | PC/EC | Кворум важнее доступности при разделе. |
| VoltDB/H-Store | PC/EC | Сильная координация в пределах конфигурации. |
Линейная согласованность означает, что все видят одну последнюю версию в реальном времени. Последовательная согласованность сохраняет один общий порядок, но он может отставать от часов. Причинная сохраняет отношения «сначала причина, потом следствие». Eventual consistency обещает, что при прекращении записей копии со временем сойдутся.
Линейность — все смотрят на один витринный манекен и сразу видят замену куртки. Eventual consistency — у магазинов несколько витрин: через минуту все обновятся, но сейчас одна ещё показывает старый размер.
Линейность дорога: нужен кворум, нельзя бездумно читать кэш, а задержка зависит от сети. Не просите её там, где пользователю всё равно, что карточка товара обновится через секунду.
Wall-clock показывает календарное время и может прыгнуть из-за NTP; monotonic clock только растёт и подходит для измерения интервалов. LWW на wall-clock опасен. Логические часы Лампорта считают порядок событий, а векторные часы показывают, какие версии уже были увидены каждой репликой.
Консенсус нужен, когда узлы должны выбрать одно значение или порядок и пережить отказ части участников. Его свойства: agreement — все решившие согласны; integrity — решение одно и предложено участником; validity — принято только допустимое значение.
FLP-невозможность говорит: в полностью асинхронной системе с одним «зависшим» узлом нельзя гарантировать завершение консенсуса за конечное время. Узлы не знают, медленный сосед или мёртвый. Поэтому Raft и Paxos используют таймауты и подозрения: обычно быстро, но безопасность сохраняется даже при ошибочном подозрении.
Кворум — большинство n/2+1. Пять узлов переживают максимум два отказа; три — один. Нечётное число даёт больше отказоустойчивости на узел стоимости: четыре не переживают больше отказов, чем три, потому что всё равно нужно три. При разделении 2+3 меньшая сторона не может безопасно решить.
Не пишите свой консенсус. Редкие сочетания задержек, GC-пауз и восстановления требуют лет тестирования. Используйте зрелый сервис и его документацию.
Lease имеет TTL. Если процесс остановился на длинной GC-паузе, lease уже истёк, но процесс после паузы может продолжить запись. Поэтому каждому владельцу выдают возрастающий fencing token, а хранилище принимает только больший токен. Таймер без fencing — не распределённая блокировка.
2PC решает атомарность конкретной транзакции через координатора и может зависнуть при его смерти. Консенсус выбирает устойчивое общее решение и порядок; это разные задачи, хотя обе требуют сетевой координации.
Брокеры · пакетная и потоковая обработка · CDC · Lambda · Kappa.
Индекс, кэш, витрина, поисковый индекс и ML-фича — это копии данных. Как только копия появилась, возникает синхронизация, а честный ответ почти всегда такой: она немного отстаёт. Проектируйте это отставание как измеримую величину и оставляйте путь к полной пересборке.
Брокер принимает сообщение и доставляет его позже. Это развязывает сервисы по времени: заказ можно принять, даже если расчёт рекомендаций временно недоступен. Брокер также сглаживает пик, позволяет отдельно масштабировать потребителей и отправлять одно событие нескольким подписчикам (fan-out).
| Стиль | Как работает | Сильная сторона | Ограничение | Выбор |
|---|---|---|---|---|
| Очередь (RabbitMQ-подобная) | Сообщение после подтверждения (ack) удаляется. | Маршрутизация, низкая задержка, небольшие задачи. | Историю нельзя перечитать и восстановить состояние. | Задачи, команды, интеграции. |
Журнал (Kafka, Kinesis, Pulsar) | Append-only записи хранятся по политике retention; потребитель хранит offset. | Replay, независимые потребители, большие объёмы. | Нужно управлять партициями, диском и схемами. | События, CDC, производные витрины. |
В очереди обычно нет возможности восстановить состояние из истории: после ack запись исчезает. В журнале запись остаётся до окончания retention, поэтому новый потребитель может начать с нужного offset и заново построить индекс.
| Семантика | Падение | Цена |
|---|---|---|
| At-most-once | Сообщение может потеряться. | Минимальная задержка и сложность. |
| At-least-once | После сбоя сообщение придёт повторно. | Нужна идемпотентность и дедупликация. |
| Exactly-once | Повторная доставка возможна, но эффект засчитан один раз. | Транзакции или идемпотентный вывод; дороже и уже. |
Exactly-once не означает «физически доставили один раз». Это гарантия эффекта: например, offset и запись результата фиксируются атомарно либо повторный результат не меняет итог.
Порядок гарантирован только внутри одной партиции. Ключ события выбирает партицию: ключ user_id сохранит порядок действий одного пользователя. Больше партиций дают параллелизм, но глобальный порядок исчезает.
Для временных ошибок используйте повторные темы (retry topics) или очереди с задержкой, экспоненциальный backoff и ограниченное число попыток. Сообщение, которое стабильно ломает обработчик, — poison message: отправьте его в dead-letter queue, сохраните причину и не блокируйте весь поток.
Для одноразовых задач выбирайте очередь. Для событий, CDC и производных данных — журнал с retention, ключом партиционирования и явной политикой replay.
MapReduce разбивает расчёт на понятные этапы: прочитать вход (read), преобразовать каждую запись (map), сгруппировать одинаковые ключи (shuffle/sort), собрать группу (reduce) и записать результат.
Сортировка по ключу — сердце подхода. После неё все записи клиента лежат рядом, поэтому соединение двух наборов можно сделать как слияние двух отсортированных потоков. Цена — сетевой обмен и запись промежуточных данных.
Пакет легко перезапустить, если вход неизменяем, расчёт детерминирован, а запись результата идемпотентна. Практический шаблон: считать партицию за день и атомарно перезаписывать date=2025-03-08, а не добавлять строки поверх старого результата.
Нужно посчитать заказы за 90 дней. Пакет читает 90 партиций, группирует по магазину и создаёт 90 дневных витрин. Сбой на 63-м дне не требует начинать всё сначала: пересчитывается только незавершённая партиция.
Batch дешевле на большом объёме, но результат может быть «вчерашним». Разбивайте наборы на партиции, планируйте backfill (пересчёт прошлого периода) и заранее считайте цену переписывания всей истории. Типичные инструменты: Spark, Hive, Flink в пакетном режиме, ClickHouse materialized views.
Поток можно представить как непрерывный batch с маленькими окнами. Важны три времени: event time — когда событие произошло, processing time — когда его обработал консьюмер, и ingest time — когда оно попало в платформу.
Окно по processing time искажает статистику при сетевом лаге или рестарте. Для продаж за час обычно правильнее event time: событие, пришедшее через две минуты, относится к часу, в котором покупка произошла.
watermark
Маркер времени: система считает, что события до этой отметки уже почти все пришли, и закрывает старые окна.
| Окно | Пример | Особенность |
|---|---|---|
| Tumbling | 00:00–01:00, 01:00–02:00 | Окна не пересекаются. |
| Hopping | 10 минут, шаг 5 минут | Окна пересекаются, обновление чаще. |
| Sliding | Последние 60 минут на каждый запрос | Плавный результат, больше состояния. |
| Session | Сессия до 30 минут бездействия | Границы зависят от активности. |
Stream–stream join требует окна: например, сопоставить клик и оплату одного заказа за 24 часа. Stream–table join обогащает событие справочником, но нужно решить, какую версию справочника использовать. Table–table поддерживает производное представление из двух меняющихся таблиц; старые версии и порядок обновлений становятся частью логики.
Состояние агрегатора хранится в state backend, часто на локальном RocksDB, с TTL и периодическими чекпоинтами. После сбоя оно восстанавливается из снапшота и replay лога. Чем больше состояние, тем дольше восстановление и тем дороже миграция. Stateful stream — один из самых сложных highload-компонентов.
Практичная гарантия — at-least-once плюс идемпотентный обработчик. Более строгий вариант — транзакционный вывод (например, Kafka transactions), где offset фиксируется вместе с результатом.
CDC (Change Data Capture) читает журнал изменений базы: PostgreSQL WAL, MySQL binlog, Mongo change streams. Сначала делается снимок, затем отправляются изменения после позиции снимка. Debezium — типичный инструмент для такого конвейера.
Материализованное представление — частный случай общего паттерна: журнал событий является источником истины, а индекс или агрегат — состоянием, которое можно вывести заново. Храните версию схемы события и измеряйте lag.
Lambda делит систему на batch layer, speed layer и serving layer. Batch пересчитывает полную историю, speed быстро добавляет свежие изменения, serving объединяет результаты. Подход оправдан, если тяжёлая историческая аналитика действительно эффективнее в пакетном движке. Главный минус — две реализации одной логики.
Kappa оставляет один поток и переигрывает журнал с начала. Это проще концептуально, но replay петабайтной истории может занять дни и дорого нагрузить хранилище. Практичный компромисс — чанкование: checkpoint плюс replay только хвоста после checkpoint.
OLTP-база заказов публикует CDC в Kafka. Потоковый агрегатор обновляет счётчики, колоночная база хранит аналитику, а кэш отдаёт фичи модели. Расхождение возможно на каждом переходе: мониторьте lag, долю устаревших фич и раз в сутки сверяйте агрегаты с источником.
Как ускорять чтение и переживать повторы, таймауты и частичные сбои.
Кэш — самая дешёвая производительность и самый дорогой баг. Очередь — самая дешёвая надёжность и источник новых требований к корректности. Ускоряйте только то, что сможете инвалидировать, проверить и восстановить.
Кэш хранит часто нужный результат ближе к потребителю. Чем ближе слой, тем меньше задержка, но тем меньше объём и тем сложнее инвалидировать данные.
| Слой | Задержка, порядок | Чем платим |
|---|---|---|
| L1 CPU | Наносекунды | Крошечный объём, сложная локальность. |
| RAM | Около 100 нс | Дороже диска, данные эфемерны. |
| Локальная сеть | 0,3–1 мс | Сеть и отдельный сервис. |
| SSD / HDD | Около 0,1 / 10 мс | IOPS и очередь запросов. |
| Межконтинентальный RTT | 100+ мс | Расстояние и нестабильность сети. |
Практические слои: локальный in-process кэш, Redis или Memcached, HTTP-кэш и CDN, буферные пулы БД и материализованные представления. Локальный кэш даёт минимум задержки, но у каждого процесса своя копия; CDN снимает трафик с приложения, но требует аккуратных заголовков приватности.
Используйте TTL, purge при записи, write-through (запись сразу в кэш и БД), write-back/write-behind (сначала кэш, позже БД), version/hash в ключе или pub/sub-инвалидацию. Самая надёжная инвалидация — сделать ключ производным от версии: product:42:v17. Старый ключ просто перестаёт использоваться.
LRU вытесняет давно не использовавшиеся записи, LFU — редко использовавшиеся. TinyLFU и W-TinyLFU лучше защищают рабочий набор от одноразового сканирования. Следите за размером рабочего набора и hit ratio как за SLO.
При hit ratio 90%, чтении кэша за 1 мс и БД за 10 мс среднее время примерно равно 0,9·1 + 0,1·10 = 1,9 мс (без учёта записи промаха). Рост hit ratio до 95% даст около 1,45 мс. Формула простая: H·Tcache + (1−H)·Tdb.
Для примера: при hit ratio 95%, чтении из кэша порядка 1 мс и из БД порядка 20 мс средняя задержка ≈ 0,95·1 + 0,05·20 = 1,95 мс против примерно 20 мс без кэша. При деградации кэша до 50% получится около 0,5·1 + 0,5·20 = 10,5 мс. Без плана прогрева кэш сам становится источником нестабильности.
Cache stampede (thundering herd) возникает, когда тысячи запросов одновременно обновляют истёкший ключ. Помогают mutex/single-flight, джиттер TTL, stale-while-revalidate и прогрев после деплоя. Горячий ключ может превысить пропускную способность одного Redis-шарда: добавьте локальный кэш второго уровня, дублируйте ключ с суффиксами или ограничьте частоту.
Защищайте кэш от отравления: проверяйте ключи, права и типы значений. Планируйте отказ Redis: сколько запросов выдержит БД при деградации, какие ответы можно отдавать устаревшими, какой rate limit включить. Иначе прогрев превратится в каскадный отказ.
Не кэшируйте баланс, остаток товара и другие данные с жёсткой консистентностью без отдельной схемы проверки. Также не стоит кэшировать редко читаемое, персонализированный ответ для каждого пользователя или данные, для которых инвалидация дороже повторного чтения.
Таймаут не равен отказу. Сервер мог списать деньги, а ответ потерялся в сети. Клиент вынужден повторить запрос; без идемпотентности это превращается в двойное списание.
| Операция | Повтор безопасен? | Как обеспечить |
|---|---|---|
GET | Да, формально. | Не менять состояние. |
PUT | Да, если задаёт итог. | Записывать одно значение по ID. |
DELETE | Да, если удаление повторно допустимо. | Считать «уже удалено» успехом. |
POST | Обычно нет. | Idempotency-Key и сохранённый ответ. |
PATCH | Зависит от операции. | Версия, CAS или абсолютное значение. |
Приёмы: клиентский Idempotency-Key с TTL и сохранением ответа, уникальное ограничение вместо «проверить, затем вставить», атомарный upsert или increment, CAS/optimistic locking с версией, дедупликация по ID сообщения. Для статусов используйте монотонную state machine: paid нельзя вернуть в pending.
Ретраи делайте с экспоненциальным backoff и джиттером, ограниченным бюджетом попыток. Таймаут внутреннего вызова должен быть меньше SLA клиента. Circuit breaker останавливает заведомо бесполезные вызовы; hedged requests применяйте осторожно и только к идемпотентным операциям.
Распределённая блокировка или lease не отменяет проблему паузы GC: владелец может пережить свой TTL. Как в §3, выдавайте возрастающий fencing token, а хранилище принимало только самый новый токен.
Клиент отправляет Idempotency-Key: pay-8f2. Сервис создаёт запись со статусом processing и уникальным ключом. Два параллельных повтора находят ту же запись: один ждёт результат, другой получает сохранённый ответ. После списания статус становится paid; повтор никогда не создаёт вторую транзакцию. При сбое после списания фоновая сверка доводит запись из processing до финального состояния.
Архитектура не заканчивается на схеме блоков: порядка 80% стоимости системы — это годы эксплуатации, а не месяцы разработки. Здесь собрано то, что делает highload-систему предсказуемой в рабочем состоянии.
Вы не можете надёжно обслуживать то, что не можете измерить. Метрики, логи и трейсы — не «для мониторинга», а единственное средство ответить на вопрос «что происходит прямо сейчас и почему», причём быстрее, чем это спросит пользователь в поддержке.
observability)Термин наблюдаемость шире, чем «мониторинг»: мониторинг проверяет известные заранее условия («диск заполнен на 90%»), а наблюдаемость — это способность задать системе новый вопрос, который вы не предвидели заранее, и получить ответ без нового релиза. Для этого нужны три разных по природе сигнала.
| Сигнал | Отвечает на вопрос | Цена | Типичный инструмент |
|---|---|---|---|
| Метрики | «Сколько, как часто, с каким трендом?» — агрегированное число во времени. | Низкая: числа компактны, дешево хранить годами. | Prometheus + Grafana, OpenTelemetry-метрики. |
| Логи | «Что именно произошло в этом одном случае?» | Высокая: текст занимает много места, растёт с трафиком линейно. | ELK-стек, Loki. |
| Трейсы | «Где в цепочке сервисов запрос потратил время?» | Средняя: нужен sampling, иначе объём как у логов. | Jaeger, Tempo, Zipkin. |
Сигналы дополняют друг друга: метрика говорит, что p99 задержки выросла в 18:03; трейс показывает, что подскочил конкретный запрос к сервису начислений; лог этого запроса — что именно там произошло (тайм-аут к внешнему API). Без всех трёх приходится либо гадать, либо чинить наугад.
RED (Rate, Errors, Duration) — набор из трёх метрик для любого сервиса, отвечающего на запросы: скорость запросов, доля ошибок, распределение длительности. Если у сервиса есть эти три графика, дежурный обычно может понять, здоров ли он, за 10 секунд.
USE (Utilization, Saturation, Errors) — то же самое для ресурса (диска, CPU, пула соединений): утилизация — доля занятого времени, насыщение (saturation) — длина очереди работы, которая не может быть обслужена немедленно, ошибки — отказы самого ресурса.
Диск занят на 95% места, но очередь операций пуста — утилизация высокая, насыщения нет, всё в порядке. Другой диск занят на 40% места, но очередь на запись — 100 операций и растёт — насыщение критическое при низкой утилизации по месту. Алертить нужно на второй случай: очередь предсказывает будущую задержку, свободное место — нет. Насыщение важнее утилизации, потому что именно оно превращается в задержку для пользователя.
Нельзя усреднять перцентили, посчитанные на разных узлах: p99 узла A и p99 узла B не складываются и не усредняются в «общий p99» — так теряется форма распределения. Правильный путь — собирать гистограммы бакетов
Значение делится на диапазоны («бакеты»): 0–5 мс, 5–10 мс, 10–25 мс и так далее. Каждый узел присылает счётчики попаданий в бакеты, а перцентиль считается уже после суммирования счётчиков со всех узлов — так распределение не искажается.
Хороший алерт — про симптом, видимый пользователю («доля ошибок 5xx выросла выше 1%»), а не про внутреннюю причину («упал один под из десяти»). Причины меняются, симптом — стабильный критерий. Page
Уведомление, которое прерывает дежурного в любое время суток (звонок, push с эскалацией) и требует немедленной реакции — в отличие от обычного тикета, который можно посмотреть утром.
Бюджет ошибок (error budget) — допустимая доля отказов за период, вытекающая из SLO. Пока бюджет не исрасходован, алерт не должен звонить: система в рамках обещания. Как только скорость расхода бюджета говорит, что он закончится раньше конца периода, — это сигнал действовать, даже если абсолютная доля ошибок пока невелика.
alert fatigue), в том числе на настоящий. Лучше меньше алертов с точным порогом, чем много «на всякий случай».Логи должны быть структурированными (JSON или другой парсимый формат, а не свободный текст), содержать уровень (debug/info/warn/error) и обязательно correlation id / trace id — идентификатор, который проходит через все сервисы одного запроса. Без него найти все записи одного инцидента среди миллионов строк почти невозможно. Держите весь трафик не на одном уровне INFO: иначе либо тонете в шуме, либо не видите проблему. Логи стоит сэмплировать (sampling) на высоком трафике — хранить не каждую запись, а представительную долю — и не хранить телефоны, паспортные данные, полные номера карт и прочий PII в открытом виде; храните ретенцию логов по политике, а не «вечно», иначе счёт за хранение обгонит счёт за вычисления.
Трейс (trace) — дерево span'ов, каждый — один шаг обработки запроса с началом, концом и родителем. Контекст трейса прокидывается через заголовки HTTP или метаданные брокера от сервиса к сервису; без этого связать span'ы разных процессов в одно дерево нельзя. Трейсинг добавляет overhead (на сериализацию и передачу контекста), поэтому в проде почти всегда включён sampling — например, 1–10% запросов целиком, плюс 100% для запросов с ошибкой. Дашборд с агрегированной задержкой скрывает единичный медленный запрос внутри среднего; трейс одного конкретного запроса, наоборот, точно показывает, какой span съел 380 из 400 мс — то есть хвост задержки трейс видит лучше любого агрегата.
Непрерывное профилирование (continuous profiling) — это тот же принцип на уровне кода: постоянный лёгкий сбор стек-трейсов CPU/памяти по всему флоту машин, чтобы можно было спросить «какая функция жрёт CPU в проде прямо сейчас» без отдельного эксперимента и без необходимости заранее знать, что искать.
trace id — инцидент приходится собирать по временным меткам вручную.INFO: критическое событие тонет среди рутинных записей о каждом запросе.Безопасность нельзя добавить патчем после релиза — она либо заложена в модель данных и границы доверия, либо становится долгом, который платится инцидентом. Начинать стоит с простой модели угроз: что именно защищаем (данные пользователей, деньги, доступность), от кого (внешний атакующий, скомпрометированный сервис, собственный сотрудник с избыточным доступом) и какова цена компрометации в деньгах и репутации — это определяет, сколько усилий оправдано.
least privilege): сервис и человек получают ровно тот доступ, который нужен для задачи, не больше.defense in depth): несколько независимых слоёв контроля, чтобы отказ одного не открывал систему целиком.at rest / in transit): диск и сеть по умолчанию читаемы кем-то ещё, если не зашифрованы.audit log): неизменяемая запись «кто, что, когда» для чувствительных операций — отдельно от обычных логов приложения.supply chain): зависимости и образы контейнеров сканируются на известные уязвимости — большая часть инцидентов приходит не из своего кода.secure by default): новый сервис без явной настройки должен быть закрыт, а не открыт всем.zero trust): mutual TLS) — оба конца соединения предъявляют сертификат и проверяют друг друга, а не только клиент проверяет сервер. Внутри периметра это заменяет неявное доверие «раз ты в нашей сети — ты свой».Производные наборы данных (аналитические выгрузки, копии для разработки, логи) должны получать персональные данные в маскированном или анонимизированном виде — полный доступ нужен единицам систем, а не каждой копии.
Честный конфликт: право на удаление (в духе GDPR) требует стереть данные конкретного человека по запросу, а append-only лог событий (см. §4) по конструкции не удаляет и не переписывает старые записи. Практичное решение — crypto shredding: персональные поля события шифруются собственным ключом на момент записи, ключ хранится отдельно от лога; «удаление» — это удаление ключа, после чего зашифрованный блок в старых событиях навсегда становится нечитаемым, а сама структура лога не нарушается.
Написать систему один раз — простая часть. Настоящая сложность highload-эксплуатации — менять её, пока она непрерывно обслуживает трафик, без простоя и без потери данных.
На таблице в миллиарды строк обычный ALTER TABLE, меняющий тип или добавляющий колонку с обязательным значением, может заблокировать запись на время, которое недопустимо для highload-сервиса, — потому что многим движкам нужно переписать весь физический файл таблицы. Стандартный безопасный приём — expand-contract (расширение → сжатие): схема меняется поэтапно, так что в каждый момент старый и новый код одновременно работают с валидной схемой.
| Шаг | Что происходит | Кто читает/пишет |
|---|---|---|
| 1. Expand | Добавляется новая колонка/таблица, ничего не удаляется. | Старый код работает как раньше, не замечая новое поле. |
| 2. Двойная запись | Код пишет одновременно в старое и новое место. | Обе версии кода видят согласованные данные. |
| 3. Backfill | Фоновый процесс с ограничением скорости (throttling) переносит старые строки в новый формат. | Не мешает обычной нагрузке на таблицу. |
| 4. Переключение чтения | Код начинает читать из нового места, проверив полноту backfill. | Новый код — основной путь. |
| 5. Contract | Старое поле/таблица перестаёт использоваться и позже удаляется. | Уже никто не пишет и не читает. |
Backfill
Заполнение нового поля или таблицы историческими данными отдельным фоновым процессом, а не одной транзакцией — обычно пачками (batch) с паузами, чтобы не создать пиковую нагрузку на движок хранения и репликацию.throttling): слишком быстрый перенос миллиардов строк создаёт ту же нагрузку, от которой expand-contract должен защищать.
Раскатка (rollout) идёт частями — сначала на часть узлов, потом на все, — значит, в переходный период одновременно работают старая и новая версия кода. Отсюда правило: новая схема должна оставаться читаемой старым кодом, а старая схема — новым кодом, пока раскатка не завершилась полностью. Это то же требование прямой/обратной совместимости, что для сериализации в §4, только на уровне таблиц, а не сообщений.
| Стратегия | Суть | Когда |
|---|---|---|
| Rolling | Узлы обновляются по одному/группами, старая и новая версия живут вместе. | Стандартный выбор для большинства сервисов. |
| Blue-green | Полный второй набор окружения, переключение трафика одним шагом. | Нужен мгновенный откат без частичного состояния. |
| Canary | Новая версия получает малую долю трафика перед полной раскаткой. | Риск обнаружить проблему до того, как она задела всех. |
| Feature flag | Код включается/выключается флагом в рантайме, без нового деплоя. | Управление поведением независимо от выпуска кода. |
Флаг — не деплой: код новой фичи может быть раскатан на все узлы, но выключен флагом для всех, кроме тестовой группы — это разделяет «доставить код» и «включить поведение», два разных решения с разным риском. Для переключения витрин данных (например, после пересборки агрегата) удобен приём partition swap: новая версия строится в отдельной партиции/таблице, а переключение — атомарная замена указателя, а не построчное обновление на живых читателях.
«Миграция вниз» (откат схемы к прошлой версии) часто физически невозможна: если contract-шаг уже удалил старую колонку и в новую записали данные, которых в старом формате не существовало, откат кода без потери данных не построить. Отсюда практика — проектировать только обратимые шаги по отдельности (expand — всегда обратим, contract — делать последним и только когда откат уже не нужен) и не совмещать в одном релизе смену схемы с рискованной логикой.
Runbook — короткий документ «что делать, если сработал этот алерт»: без него дежурный тратит время инцидента на выяснение контекста, а не на исправление. Постмортем без обвинений (blameless postmortem) разбирает инцидент как свойство системы и процесса, а не ошибку конкретного человека — это единственный способ получить честное описание того, что произошло, вместо защитной версии. Бюджет ошибок из §6.1 — не только повод для алерта, но и организационный инструмент: если он исрасходован, команда обязана приостановить новые фичи и заняться надёжностью, прежде чем продолжать риск.
Стоимость хранения — часть архитектуры, а не только финансовый вопрос: горячие данные (частый доступ) держат на быстром и дорогом хранилище, холодные — сжимают и переносят на дешёвое; TTL (время жизни записи) автоматически убирает данные, которые никому не нужны после срока. Компрессия логов и старых партиций обычно даёт кратное сокращение объёма почти без потери возможности разобрать инцидент постфактум.
deployment pipeline) — автоматические тесты, постепенная раскатка, метрики после релиза, быстрый откат — это то, что превращает «часто менять» и «стабильно работать» из противоречия в одновременно достижимую пару свойств.Если сделать одним ALTER TABLE ... RENAME COLUMN в лоб на такой таблице во многих движках: либо операция блокирует запись на время, недопустимое для продовольствия highload-сервиса, либо старый код, ещё обращающийся к прежнему имени колонки, немедленно начинает падать с ошибкой — потому что деплой кода и миграция схемы не атомарны друг относительно друга.
Через expand-contract: (1) добавляем новую колонку рядом со старой; (2) код переходит на запись в обе колонки одновременно; (3) фоновый backfill с throttling переносит 2 млрд старых значений порциями, не мешая обычному трафику; (4) после проверки полноты backfill код переключается читать из новой колонки; (5) через один-два релиза, когда весь флот точно обновлён, старая колонка перестаёт заполняться и позже удаляется отдельной низкоприоритетной миграцией.
Итог §6 Наблюдаемость превращает инцидент из «непонятно, что случилось» в «известно, где смотреть, за минуты». Безопасность превращает риск компрометации в управляемую, а не игнорируемую величину. Экспанд-контракт и постепенная раскатка превращают изменение работающей системы из рискованного события в рутинную операцию с предсказуемым откатом. Все три темы — не дополнение к архитектуре, а её продолжение во времени: то, ради чего вообще имеет смысл проектировать репликацию, партиционирование и производные данные из §2–§5 — чтобы система оставалась управляемой годы после первого релиза.
Ответы должны помещаться на одну страницу. Если на вопрос нет числового ответа — это не решение, а надежда.
ADR, Architecture Decision Record) с явными ответами на эти 15 вопросов, датой и автором — этого достаточно для старта разработки и для будущего ревью решения. По мере роста зрелости проекта в тот же процесс стоит добавить: версионирование API, сериализацию и совместимость схем сообщений, работу с часами и метками времени, ёмкостное планирование (capacity planning), регулярные учения disaster recovery (проверка восстановления из бэкапа на практике, а не только на бумаге) и оценку стоимости инфраструктуры на 1000 пользователей.