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.
Исходи из того, что частичный отказ — это норма: в любой момент какая-то нода лежит, какое-то сообщение потеряно или опоздало, а часы расходятся. Хороший ответ всегда покрывает: что может сломаться → как это обнаружить → как система ведёт себя во время сбоя → как восстанавливается → какая гарантия при этом сохраняется.
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: синхронные строгие чтения стоят задержки; быстрые чтения могут быть слегка устаревшими.
2. Consistency models (the spectrum)2. Модели согласованности (спектр)
From "always correct, slower" to "fast, but surprising".От «всегда верно, но медленнее» до «быстро, но с сюрпризами».
| Model | Guarantee | Use for |
|---|---|---|
| Linearizable (strong) | Every read returns the latest committed write; behaves like one copy. | Locks, leader election, balances, config. |
| Causal | Cause-and-effect operations seen in order by everyone; concurrent ones may differ. | Comments/replies, messaging. |
| Read-your-writes | A client always sees its own previous writes. | "I posted but don't see it" bugs. |
| Eventual | If writes stop, all replicas converge — eventually. No ordering promise. | Caches, feeds, analytics, metrics. |
| Модель | Гарантия | Где применять |
|---|---|---|
| Линеаризуемая (строгая) | Любое чтение возвращает последнюю закоммиченную запись; ведёт себя как одна копия. | Блокировки, выбор лидера, балансы, конфиги. |
| Причинная (causal) | Причинно-связанные операции все видят по порядку; конкурентные могут отличаться. | Комментарии/ответы, мессенджеры. |
| Read-your-writes | Клиент всегда видит свои предыдущие записи. | Баги «я запостил, но не вижу». |
| 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 перемешивает почти всё.
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 ключей, а не все. Виртуальные ноды (много точек кольца на машину) держат нагрузку ровной, а ребаланс — пропорциональным.
4. Replication & quorums4. Репликация и кворумы
| Topology | How | Main trade-off |
|---|---|---|
| Leader–follower | All writes to one leader, copied to followers. | Simple, no conflicts; but leader is a bottleneck/SPOF and async followers lag. |
| Multi-leader | Several 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, множества чтения и записи всегда пересекаются хотя бы на одной реплике — поэтому чтение гарантированно видит последнюю подтверждённую запись.
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 решает ту же задачу и славится тем, что в нём труднее разобраться.
6. Failure modes & mitigations6. Режимы сбоев и меры
| Failure | Detect | Mitigation |
|---|---|---|
| Node crash | Missed heartbeats / lease expiry | Replication + failover; replay from a durable log |
| Gray failure (slow node) | p99 latency, error rate vs peers | Hedged requests, timeouts, eject from pool, circuit breaker |
| Network partition | Quorum loss, heartbeat gaps | CP: minority stops; AP: degrade + reconcile; fencing tokens |
| Clock skew | NTP drift monitoring | Logical clocks; never trust wall-clock for ordering |
| Poison message | Same record retried forever | Dead-letter queue, retry cap, skip-and-log |
| Retry storm | Spike right after recovery | Exponential backoff + jitter, retry budgets |
| Cascading failure | Saturation spreading to callers | Bulkheads, timeouts everywhere, backpressure, load shedding |
| Сбой | Обнаружение | Мера |
|---|---|---|
| Падение ноды | Пропуск heartbeat / истёк lease | Репликация + failover; перечитка из надёжного лога |
| Серый сбой (медленная нода) | p99-задержка, error rate против соседей | Hedged-запросы, таймауты, вывод из пула, circuit breaker |
| Сетевой разрыв | Потеря кворума, пропуски heartbeat | CP: меньшинство останавливается; AP: деградация + сверка; fencing tokens |
| Расхождение часов | Мониторинг дрейфа NTP | Логические часы; никогда не доверять wall-clock для порядка |
| Poison message | Одна запись ретраится вечно | Dead-letter queue, лимит ретраев, skip-and-log |
| Шторм ретраев | Всплеск сразу после восстановления | Экспоненциальный backoff + jitter, бюджеты ретраев |
| Каскадный отказ | Насыщение расползается на вызывающих | Bulkheads, таймауты везде, backpressure, сброс нагрузки |
7. Delivery semantics & exactly-once7. Семантики доставки и exactly-once
| Semantic | Means | Risk |
|---|---|---|
| At-most-once | 0 or 1 delivery (no retry) | Data loss |
| At-least-once | 1+ deliveries (retry until acked) | Duplicates → double counting |
| Exactly-once (effect) | At-least-once + dedup / idempotent writes | Cost, complexity |
| Семантика | Значит | Риск |
|---|---|---|
| At-most-once | 0 или 1 доставка (без ретраев) | Потеря данных |
| At-least-once | 1+ доставок (ретрай до ack) | Дубли → двойной счёт |
| Exactly-once (эффект) | At-least-once + дедуп / идемпотентная запись | Цена, сложность |
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.
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 — сигнал обратной связи, который тормозит продюсера, когда консьюмер не успевает — вместо того чтобы буферы росли до нехватки памяти. Без него короткое замедление превращается в падение.
- 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): гасит всплески; лог — амортизатор.
9. Batch vs streaming & watermarks9. Batch vs streaming и watermarks
| Batch | Streaming | |
|---|---|---|
| Latency | Minutes–hours | ms–seconds |
| Correctness | Easy (complete data, reprocess freely) | Harder (late / out-of-order data) |
| Use for | Daily aggregates, backfills, training data | Alerting, fraud, real-time dashboards |
| Batch | Streaming | |
|---|---|---|
| Задержка | Минуты–часы | мс–секунды |
| Корректность | Легко (полные данные, свободная переобработка) | Труднее (опоздавшие / неупорядоченные данные) |
| Для чего | Дневные агрегаты, бэкфилы, обучающие данные | Алертинг, фрод, дашборды в реальном времени |
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 или проход батч-сверки, чтобы уточнить число.
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.Ответь вслух, потом переверни.