Distributed Systems — the core ideasРаспределённые системы — ключевые идеи

The real names and the working knowledge: CAP/PACELC, consistency, replication & quorums, consensus, failures, exactly-once, backpressure, watermarks. Comfortable middle level — done with this, the Advanced page is just more depth and edge cases. Настоящие термины и рабочее понимание: CAP/PACELC, согласованность, репликация и кворумы, консенсус, сбои, exactly-once, backpressure, watermarks. Комфортный средний уровень — после него «Продвинутый» это лишь больше глубины и краевых случаев.
🟢 BeginnerНовичок 🟡 IntermediateСредний 🔴 AdvancedПродвинутый

0. The one mental model0. Одна мысленная модель

Assume partial failure is normal: at any moment some node is down, some message is lost or late, and clocks disagree. A good answer always covers: what can fail → how you detect it → how the system behaves during the failure → how it recovers → what guarantee still holds.

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

Keep this anchorДержи этот якорь You can never tell "dead" from "slow". So make every write retry-safe (idempotent), every consumer replay-safe (dedup + checkpoint), and make timeouts/backpressure explicit so a slow dependency degrades instead of collapsing. Никогда нельзя отличить «мёртв» от «тормозит». Поэтому делай каждую запись безопасной к повтору (идемпотентной), каждого консьюмера — безопасным к перечитке (дедуп + чекпоинт), а таймауты/backpressure — явными, чтобы медленная зависимость деградировала, а не обрушивалась.

1. CAP and PACELC1. CAP и PACELC

CAP: during a network Partition, a system can keep at most one of Consistency (every read sees the latest write) or Availability (every node still answers). When there's no partition, you can have both — CAP only describes the partition case.

CAP: во время сетевого разрыва (Partition) система может сохранить максимум одно из Consistency (любое чтение видит последнюю запись) или Availability (каждая нода всё ещё отвечает). Когда разрыва нет — можно и то, и другое; CAP описывает только случай разрыва.

  • CP refuse on the cut-off side to stay correct (e.g. ZooKeeper, etcd, HBase).
  • AP keep answering, reconcile later, tolerate stale reads (e.g. Cassandra, Dynamo).
  • "CA" isn't a real choice — partitions will happen, so you only choose C-vs-A for the partition.
  • CP отказывать на отрезанной стороне, чтобы остаться верным (например ZooKeeper, etcd, HBase).
  • AP продолжать отвечать, сверяться позже, терпеть устаревшие чтения (например Cassandra, Dynamo).
  • «CA» — не реальный выбор: разрывы будут, поэтому C-или-A выбирают только на время разрыва.

PACELC is the more useful version: if Partition → trade A vs C; Else (healthy) → trade Latency vs Consistency. Since there's almost never a partition, the everyday trade-off is the EL/EC one: synchronous strong reads cost latency; fast reads may be slightly stale.

PACELC — более полезная версия: если Partition → выбор A vs C; иначе (Else) при здоровой сети → выбор Latency vs Consistency. Так как разрыва почти никогда нет, ежедневный компромисс — именно EL/EC: синхронные строгие чтения стоят задержки; быстрые чтения могут быть слегка устаревшими.

Common gotchaЧастая ловушка "Consistency" in CAP means linearizability (a very strong model), NOT the "C" in ACID. And labels aren't binary — many systems are tunable (e.g. Cassandra QUORUM). Talk about knobs, not fixed labels. «Consistency» в CAP — это линеаризуемость (очень строгая модель), а НЕ «C» из ACID. И ярлыки не бинарны — многие системы настраиваемы (например Cassandra с QUORUM). Говори про «ручки» настройки, а не про фиксированные ярлыки.

2. Consistency models (the spectrum)2. Модели согласованности (спектр)

From "always correct, slower" to "fast, but surprising".От «всегда верно, но медленнее» до «быстро, но с сюрпризами».

STRONGER / SLOWER ───────────────────────▶ WEAKER / FASTER Linearizable ─ Causal ─ Read-your-writes ─ Eventual
ModelGuaranteeUse for
Linearizable (strong)Every read returns the latest committed write; behaves like one copy.Locks, leader election, balances, config.
CausalCause-and-effect operations seen in order by everyone; concurrent ones may differ.Comments/replies, messaging.
Read-your-writesA client always sees its own previous writes."I posted but don't see it" bugs.
EventualIf writes stop, all replicas converge — eventually. No ordering promise.Caches, feeds, analytics, metrics.
МодельГарантияГде применять
Линеаризуемая (строгая)Любое чтение возвращает последнюю закоммиченную запись; ведёт себя как одна копия.Блокировки, выбор лидера, балансы, конфиги.
Причинная (causal)Причинно-связанные операции все видят по порядку; конкурентные могут отличаться.Комментарии/ответы, мессенджеры.
Read-your-writesКлиент всегда видит свои предыдущие записи.Баги «я запостил, но не вижу».
Eventual (итоговая)Если записи прекратились, все реплики сойдутся — со временем. Порядок не гарантирован.Кэши, ленты, аналитика, метрики.
"Eventual" is underspecified«Eventual» — недосказана Plain eventual gives no session guarantees: a user can write, then read an older replica and not see it. Read-your-writes and causal are the practical session guarantees that make eventual systems usable for humans. Чистая eventual не даёт сессионных гарантий: пользователь может записать, затем прочитать старую реплику и не увидеть запись. Read-your-writes и causal — практические сессионные гарантии, которые делают eventual-системы пригодными для людей.

3. Partitioning / sharding3. Партиционирование / шардирование

How you split data decides hotspots and how painful it is to add a node.Как делишь данные — так и получаешь горячие точки и боль при добавлении ноды.

Hash partitioningХеш-партиционирование

node = hash(key) % N. Even spread, kills hotspots — but no range queries, and changing N reshuffles almost everything.

node = hash(key) % N. Ровное распределение, убирает горячие точки — но нет диапазонных запросов, и смена N перемешивает почти всё.

even loadровноno rangesнет диапазонов

Range partitioningДиапазонное партиционирование

Contiguous ranges per shard (e.g. by date). Great for scans, but sequential keys create hotspots (the "always newest" trap).

Непрерывные диапазоны на шард (например по дате). Отлично для сканов, но последовательные ключи дают горячие точки (ловушка «всё самое новое»).

range scansсканы диапазоновhot rangeгорячий диапазон

Consistent hashing — how you add a node cheaplyConsistent hashing — как дёшево добавить ноду

Keys and nodes both hash onto a ring; a key belongs to the next node clockwise. Adding/removing a node only moves about 1/N of keys, not all of them. Virtual nodes (many ring points per machine) keep the load smooth and rebalancing proportional.

Ключи и ноды хешируются на кольцо; ключ принадлежит следующей ноде по часовой. Добавление/удаление ноды двигает лишь около 1/N ключей, а не все. Виртуальные ноды (много точек кольца на машину) держат нагрузку ровной, а ребаланс — пропорциональным.

Hotspot fix (ties to Spark)Лечение горячей точки (связь со Spark) A hot shard usually means a skewed key (a "celebrity"/whale, a null, or a sequential key). Fixes: salt the key (add a bucket suffix), use a composite key, or split the hot range. This is the exact same cure as a skewed join key in Spark. Горячий шард обычно значит перекошенный ключ («звезда»/кит, null или последовательный ключ). Лечение: солить ключ (добавить суффикс-бакет), составной ключ или разбить горячий диапазон. Это ровно то же лечение, что и перекошенный ключ джоина в Spark.

4. Replication & quorums4. Репликация и кворумы

TopologyHowMain trade-off
Leader–followerAll writes to one leader, copied to followers.Simple, no conflicts; but leader is a bottleneck/SPOF and async followers lag.
Multi-leaderSeveral leaders accept writes (e.g. per region).Fast local writes; but write conflicts need resolving.
Leaderless (quorum)Client writes/reads many replicas; a quorum decides.Highly available, tunable; you reason about quorums yourself.
ТопологияКакГлавный компромисс
Лидер–фолловерВсе записи в одного лидера, копируются фолловерам.Просто, нет конфликтов; но лидер — узкое место/SPOF, а асинхронные фолловеры отстают.
Мульти-лидерНесколько лидеров принимают записи (например по регионам).Быстрые локальные записи; но конфликты записи надо разрешать.
Без лидера (кворум)Клиент пишет/читает много реплик; решает кворум.Высокая доступность, настраиваемо; кворумы считаешь сам.

The quorum rule: R + W > NПравило кворума: R + W > N

With N replicas, read from R and write to W. If R + W > N, the read set and write set always overlap by at least one replica — so a read is guaranteed to see the latest acknowledged write.

При N репликах читаем из R и пишем в W. Если R + W > N, множества чтения и записи всегда пересекаются хотя бы на одной реплике — поэтому чтение гарантированно видит последнюю подтверждённую запись.

N = 3 W=2, R=2 → R+W = 4 > 3 ✓ overlap (balanced) W=1, R=1 → R+W = 2 ≤ 3 ✗ may read stale (pure eventual)
Sync vs asyncSync vs async Sync replication = leader waits for a follower's ack → no data loss on crash, but slower. Async = fast, but a crash can lose the un-replicated tail. Semi-sync (wait for ≥1 follower) is the common middle ground. Синхронная репликация = лидер ждёт ack фолловера → нет потерь при падении, но медленнее. Асинхронная = быстро, но падение может потерять нереплицированный «хвост». Semi-sync (ждать ≥1 фолловера) — частая золотая середина.

5. Consensus & leader election5. Консенсус и выбор лидера

Consensus = getting a group of nodes to agree on one value (the next log entry, or who is leader) despite crashes and delays. It needs a strict majority to commit, which is why these clusters are odd-sized (3, 5): a partition leaves at most one side with a majority, so only that side makes progress — no split-brain. 3 nodes tolerate 1 failure; 5 tolerate 2.

Консенсус = заставить группу нод согласиться об одном значении (следующая запись лога или кто лидер) несмотря на падения и задержки. Чтобы закоммитить, нужно строгое большинство — поэтому такие кластеры нечётны (3, 5): при разрыве большинство будет максимум у одной стороны, только она движется вперёд — нет split-brain. 3 ноды переживают 1 отказ; 5 — 2 отказа.

Raft is the teachable algorithm: a randomized timeout triggers an election, a candidate that wins a majority becomes leader for a term, and a log entry commits once a majority has it. Paxos solves the same problem and is famously harder to reason about.

Raft — «понятный» алгоритм: случайный таймаут запускает выборы, кандидат, набравший большинство, становится лидером на term, а запись лога коммитится, когда она есть у большинства. Paxos решает ту же задачу и славится тем, что в нём труднее разобраться.

Fencing tokens (says senior)Fencing tokens (признак сеньора) Electing a leader isn't enough. A "slow-then-revived" old leader must be stopped: give each leader a monotonic fencing token (an ever-increasing id), and have storage reject any write carrying an older token. This prevents the #1 cause of split-brain corruption. Выбрать лидера мало. «Затормозившего, но ожившего» старого лидера надо остановить: дай каждому лидеру монотонный fencing token (всё возрастающий id), и пусть хранилище отвергает любую запись со старым токеном. Это предотвращает главную причину порчи данных при split-brain.

6. Failure modes & mitigations6. Режимы сбоев и меры

FailureDetectMitigation
Node crashMissed heartbeats / lease expiryReplication + failover; replay from a durable log
Gray failure (slow node)p99 latency, error rate vs peersHedged requests, timeouts, eject from pool, circuit breaker
Network partitionQuorum loss, heartbeat gapsCP: minority stops; AP: degrade + reconcile; fencing tokens
Clock skewNTP drift monitoringLogical clocks; never trust wall-clock for ordering
Poison messageSame record retried foreverDead-letter queue, retry cap, skip-and-log
Retry stormSpike right after recoveryExponential backoff + jitter, retry budgets
Cascading failureSaturation spreading to callersBulkheads, timeouts everywhere, backpressure, load shedding
СбойОбнаружениеМера
Падение нодыПропуск heartbeat / истёк leaseРепликация + failover; перечитка из надёжного лога
Серый сбой (медленная нода)p99-задержка, error rate против соседейHedged-запросы, таймауты, вывод из пула, circuit breaker
Сетевой разрывПотеря кворума, пропуски heartbeatCP: меньшинство останавливается; AP: деградация + сверка; fencing tokens
Расхождение часовМониторинг дрейфа NTPЛогические часы; никогда не доверять wall-clock для порядка
Poison messageОдна запись ретраится вечноDead-letter queue, лимит ретраев, skip-and-log
Шторм ретраевВсплеск сразу после восстановленияЭкспоненциальный backoff + jitter, бюджеты ретраев
Каскадный отказНасыщение расползается на вызывающихBulkheads, таймауты везде, backpressure, сброс нагрузки
Gray failure is the killerСерый сбой — главный убийца A crash is easy: the node is gone, failover fires. A gray failure (30% slow, dropping 2% of requests, GC-pausing) passes naive health checks while poisoning tail latency — the cause of most real incidents. Detect with differential metrics (this replica vs peers) and hedged requests. Падение — это легко: ноды нет, срабатывает failover. Серый сбой (на 30% медленнее, теряет 2% запросов, встаёт на GC) проходит наивные health-check'и, отравляя tail latency — причина большинства реальных инцидентов. Лови дифференциальными метриками (эта реплика vs соседи) и hedged-запросами.

7. Delivery semantics & exactly-once7. Семантики доставки и exactly-once

SemanticMeansRisk
At-most-once0 or 1 delivery (no retry)Data loss
At-least-once1+ deliveries (retry until acked)Duplicates → double counting
Exactly-once (effect)At-least-once + dedup / idempotent writesCost, complexity
СемантикаЗначитРиск
At-most-once0 или 1 доставка (без ретраев)Потеря данных
At-least-once1+ доставок (ретрай до ack)Дубли → двойной счёт
Exactly-once (эффект)At-least-once + дедуп / идемпотентная записьЦена, сложность
The truth about exactly-onceПравда об exactly-once Exactly-once delivery over a network is impossible (you can never confirm your ack arrived — the Two Generals problem). What you achieve is exactly-once effect: accept duplicate delivery, but make the effect idempotent. Formula: at-least-once + idempotent write (upsert by key) + dedup (idempotency key / offset) + atomic commit of write&offset. Exactly-once доставка по сети невозможна (нельзя подтвердить, что твой ack дошёл — проблема двух генералов). Достижим exactly-once эффект: принимай дубли доставки, но делай эффект идемпотентным. Формула: at-least-once + идемпотентная запись (upsert по ключу) + дедуп (idempotency key / оффсет) + атомарный коммит записи и оффсета.

In practice: Kafka EOS = idempotent producer + transactions + read_committed (exactly-once within Kafka; an external sink still needs to be idempotent). Flink = checkpoints + two-phase-commit sink. Spark Structured Streaming = checkpointed offsets + idempotent MERGE upsert.

На практике: Kafka EOS = идемпотентный продюсер + транзакции + read_committed (exactly-once внутри Kafka; внешнему приёмнику всё равно нужна идемпотентность). Flink = чекпоинты + sink с two-phase-commit. Spark Structured Streaming = чекпоинт оффсетов + идемпотентный MERGE-upsert.

Design ruleПравило дизайна Prefer naturally idempotent writes: upsert by key and SET x = v (not x = x + 1). If you must aggregate, dedup first on a stable event_id so reprocessing is safe. Предпочитай естественно идемпотентные записи: upsert по ключу и SET x = v (а не x = x + 1). Если надо агрегировать, сначала дедуп по стабильному event_id, чтобы переобработка была безопасной.

8. Backpressure & flow control8. Backpressure и контроль потока

Backpressure is a feedback signal that slows the producer when the consumer can't keep up — instead of letting buffers grow until the process runs out of memory. Without it, a brief slowdown becomes a crash.

Backpressure — сигнал обратной связи, который тормозит продюсера, когда консьюмер не успевает — вместо того чтобы буферы росли до нехватки памяти. Без него короткое замедление превращается в падение.

FAST PRODUCER ──▶ [ bounded queue ] ──▶ SLOW CONSUMER │ full? ▼ signal upstream: slow down / shed
  • Pull-based / demand (Reactive Streams, Flink credit-based): consumer asks for N items; producer can't exceed demand.
  • Bounded buffer + blocking: producer blocks when full (TCP flow control works like this).
  • Load shedding: drop low-priority work when overloaded — better to drop 1% than fall over.
  • Buffer to a durable log (Kafka): absorbs spikes; the log is the shock absorber.
  • Pull / по требованию (Reactive Streams, кредитный Flink): консьюмер просит N элементов; продюсер не превышает спрос.
  • Ограниченный буфер + блокировка: продюсер блокируется при заполнении (так работает TCP flow control).
  • Сброс нагрузки: отбрасывать низкоприоритетное при перегрузе — лучше отбросить 1%, чем упасть.
  • Буфер в надёжный лог (Kafka): гасит всплески; лог — амортизатор.
Kafka buys time, not infinityKafka даёт время, не бесконечность A log absorbs spikes, but if the consumer is chronically slower, lag grows until retention deletes un-consumed data — silent loss. Always monitor consumer lag and scale partitions/consumers. Лог гасит всплески, но если консьюмер хронически медленнее, лаг растёт, пока retention не удалит непрочитанное — тихая потеря. Всегда мониторь лаг консьюмера и масштабируй партиции/консьюмеров.

9. Batch vs streaming & watermarks9. Batch vs streaming и watermarks

BatchStreaming
LatencyMinutes–hoursms–seconds
CorrectnessEasy (complete data, reprocess freely)Harder (late / out-of-order data)
Use forDaily aggregates, backfills, training dataAlerting, fraud, real-time dashboards
BatchStreaming
ЗадержкаМинуты–часымс–секунды
КорректностьЛегко (полные данные, свободная переобработка)Труднее (опоздавшие / неупорядоченные данные)
Для чегоДневные агрегаты, бэкфилы, обучающие данныеАлертинг, фрод, дашборды в реальном времени

Event time, processing time, watermarksEvent time, processing time, watermarks

Events arrive late and out of order. You usually aggregate by event time (when it happened), not processing time (when you saw it). A watermark is the system's claim "I've now seen all events up to time W" — it decides when a time window is allowed to close.

События приходят с опозданием и вне порядка. Обычно агрегируешь по event time (когда случилось), а не по processing time (когда увидел). Watermark — утверждение системы «я уже видел все события до момента W» — оно решает, когда временно́му окну можно закрыться.

  • Tighter watermark → lower latency, but more late data dropped.
  • Looser watermark → more complete, but higher latency and more state held.
  • Very-late events → side-output / dead-letter, or a batch reconciliation pass to true-up the number.
  • Туже watermark → меньше задержка, но больше опоздавших отбрасывается.
  • Свободнее watermark → полнее, но выше задержка и больше удерживаемого состояния.
  • Сильно опоздавшие → side-output / dead-letter или проход батч-сверки, чтобы уточнить число.
Don't order by wall-clockНе упорядочивай по wall-clock Processing-time windows give wrong results whenever ingestion lags or you replay. Key correctness on a trustworthy event_time in the payload; treat the watermark as your completeness-vs-latency dial. Окна по processing-time дают неверный результат при отставании ingestion или перечитке. Опирай корректность на надёжный event_time в payload; watermark — это «ручка» полнота-vs-задержка.

10. Self-quiz10. Самопроверка

Answer out loud, then flip.Ответь вслух, потом переверни.

What does PACELC add over CAP?Что PACELC добавляет к CAP?
tapнажми
The Else case: when healthy (almost always), you still trade Latency vs Consistency. CAP only covers behavior during a partition; PACELC covers the everyday path too.Случай Else: при здоровой сети (почти всегда) всё равно есть компромисс Latency vs Consistency. CAP описывает только поведение при разрыве; PACELC — ещё и обычный путь.
N=5: pick R, W to survive 2 write-path failures with overlap.N=5: подбери R, W, чтобы пережить 2 отказа на записи с пересечением.
tapнажми
W=3, R=3 → R+W=6 > 5 (overlap) and W=3 is a majority (tolerates 2 of 5 down). Note: quorum overlap is not full linearizability under concurrent writes.W=3, R=3 → R+W=6 > 5 (пересечение), и W=3 — большинство (переживает 2 из 5). Заметь: пересечение кворума ≠ полная линеаризуемость при конкурентных записях.
Is exactly-once delivery possible?Возможна ли exactly-once доставка?
tapнажми
No (Two Generals). You get exactly-once effect = at-least-once + idempotent write + dedup + atomic commit of write & offset.Нет (два генерала). Достигаешь exactly-once эффект = at-least-once + идемпотентная запись + дедуп + атомарный коммит записи и оффсета.
Why odd cluster sizes, and what stops a revived old leader?Почему нечётный размер кластера и что останавливает ожившего старого лидера?
tapнажми
Only one side of a partition can hold a majority → no split-brain. A revived stale leader is stopped by fencing tokens: storage rejects writes carrying an older token.Большинство может быть только у одной стороны разрыва → нет split-brain. Ожившего старого лидера останавливают fencing tokens: хранилище отвергает записи со старым токеном.
Why is gray failure harder than a crash?Почему серый сбой труднее падения?
tapнажми
A crash trips health checks and fails over. A gray failure stays "up" while slow/partially-erroring — it passes naive checks and poisons tail latency. Detect with differential metrics + hedged requests.Падение валит health-check и делает failover. Серый сбой остаётся «живым», будучи медленным/частично ошибочным — проходит наивные проверки и отравляет tail latency. Лови дифференциальными метриками + hedged-запросами.
What does tuning a watermark cost?Чем стоит настройка watermark?
tapнажми
Tighter → lower latency but more dropped late data. Looser → more complete but higher latency + more state. Very-late data → side-output or a batch true-up.Туже → меньше задержка, но больше отброшенных опоздавших. Свободнее → полнее, но выше задержка + больше состояния. Сильно опоздавшие → side-output или батч-уточнение.