Distributed Systems (Advanced) — DE Interview PrepРаспределённые системы (Продвинутый) — Подготовка к DE-интервью

Fast, dense revision for Senior/Lead Data Engineering interviews. Designed for engineers who reason about failure modes and optimize billion-event pipelines. Frame answers with your real stack: Spark-on-Hadoop / Scala, Kafka, Airflow, and billions of events/day. Быстрое, плотное повторение для интервью на Senior/Lead Data Engineering. Рассчитано на инженеров, которые рассуждают о режимах сбоев и оптимизируют пайплайны с миллиардами событий. Обрамляй ответы своим реальным стеком: Spark-on-Hadoop / Scala, Kafka, Airflow и миллиарды событий/день.
CAP · PACELC Quorums R+W>N Raft / Paxos Exactly-once (really)Exactly-once (правда) Backpressure Data skew & shufflesПерекос данных и shuffle Watermarks & late dataWatermark и поздние данные
🟢 BeginnerНовичок 🟡 IntermediateСредний 🔴 AdvancedПродвинутый

00 Mental model: how to reason about failure modesМентальная модель: как рассуждать о режимах сбоев

The single thread that runs through every distributed-systems answer. Lead with this framing.

Единая нить, проходящая через каждый ответ о распределённых системах. Начинай с этого фрейминга.

Distributed systems is the study of what happens when independent failure meets unbounded message delay. A senior answer always names: (1) what can fail, (2) how you detect it, (3) how the system behaves during the failure, (4) how it recovers, and (5) what guarantee survives the whole thing.

Распределённые системы — это изучение того, что происходит, когда независимый сбой встречается с неограниченной задержкой сообщений. Ответ senior-уровня всегда называет: (1) что может упасть, (2) как ты это обнаруживаешь, (3) как система ведёт себя во время сбоя, (4) как она восстанавливается, и (5) какая гарантия переживает всё это.

The reasoning checklist (say it out loud)

Чек-лист рассуждений (произноси вслух)

  1. Enumerate failure domains: node crash, slow node (gray failure), disk full, network partition, clock skew, dependency down, poison message, traffic spike.
  2. Pick the model: crash-stop vs crash-recovery vs Byzantine; synchronous vs partially-synchronous network (real world = partially synchronous).
  3. Detection: timeouts and heartbeats are suspicion, not truth. You can never distinguish "dead" from "slow" — this is the FLP impossibility in practice.
  4. Blast radius: does one failure cascade? Bulkheads, timeouts, circuit breakers, backpressure.
  5. Recovery + guarantee: idempotency, retries with jitter, replay from durable log, checkpoint/offset rewind.
  1. Перечисли домены сбоев: падение ноды, медленная нода (gray failure), диск полон, разделение сети, перекос часов, зависимость упала, poison-сообщение, всплеск трафика.
  2. Выбери модель: crash-stop vs crash-recovery vs Byzantine; синхронная vs частично синхронная сеть (реальный мир = частично синхронная).
  3. Обнаружение: таймауты и heartbeat — это подозрение, не истина. Ты никогда не различишь «мёртвый» от «медленный» — это FLP impossibility на практике.
  4. Радиус поражения: один сбой каскадирует? Bulkheads, таймауты, circuit breakers, backpressure.
  5. Восстановление + гарантия: идемпотентность, ретраи с jitter, переигрывание из надёжного лога, откат checkpoint/offset.
Senior framing "I assume partial failure is the steady state, not an exception. I design every write to be retry-safe (idempotent), every consumer to be replay-safe (dedup + checkpoint), and I make timeouts/backpressure explicit so a slow dependency degrades instead of collapsing."
Фрейминг senior-уровня «Я исхожу из того, что частичный сбой — это норма, а не исключение. Я делаю каждую запись безопасной для ретраев (идемпотентной), каждого консьюмера безопасным для переигрывания (дедуп + checkpoint), и делаю таймауты/backpressure явными, чтобы медленная зависимость деградировала, а не рушилась.»
Interviewer will probe "A node stops responding. Is it dead or slow, and does it matter?" — Correct answer: you cannot tell them apart, so the protocol must be correct under either. Fencing tokens / leases stop a "slow-then-revived" leader from doing damage.
О чём спросит интервьюер «Нода перестала отвечать. Она мёртвая или медленная, и важно ли это?» — Правильный ответ: ты не можешь их различить, поэтому протокол должен быть корректен при любом. Fencing tokens / leases останавливают «медленно-потом-ожившего» лидера от нанесения вреда.

01 CAP theorem & PACELCТеорема CAP и PACELC

The most-asked opener. Be precise: CAP is about behavior during a partition, not a permanent property.

Самый частый вопрос на открытие. Будь точен: CAP — про поведение во время разделения, а не постоянное свойство.

CAP, stated correctly

CAP, сформулированная правильно

During a network Partition, a distributed store can preserve at most one of Consistency (linearizable reads) or Availability (every non-failing node answers). When there is no partition, you get both — CAP says nothing about the happy path.

Во время разделения сети (Partition) распределённое хранилище может сохранить максимум одно из Consistency (линеаризуемое чтение) или Availability (каждая не-упавшая нода отвечает). Когда разделения нет, получаешь оба — CAP ничего не говорит про happy path.

  • CP — refuse/block on the minority side to stay consistent. (ZooKeeper, etcd, HBase, Spanner-as-CP.)
  • AP — keep answering, reconcile later, tolerate stale/conflicting reads. (Dynamo, Cassandra default, Riak.)
  • "CA" is not a real operating point for a networked system — partitions will happen, so you only choose C-vs-A for the partition.
  • CP — отказывать/блокироваться на стороне меньшинства, чтобы сохранить консистентность. (ZooKeeper, etcd, HBase, Spanner-как-CP.)
  • AP — продолжать отвечать, согласовывать потом, терпеть устаревшие/конфликтующие чтения. (Dynamo, Cassandra по умолчанию, Riak.)
  • «CA» — не реальная точка работы для сетевой системы — разделения будут происходить, поэтому выбираешь только C-против-A для разделения.
Common gotcha "Consistency" in CAP = linearizability, a very strong model. It is NOT the "C" in ACID. Many "AP" systems still give per-key strong-ish guarantees; many "CP" systems are tunable. Avoid binary labels; talk about knobs.
Частая ошибка «Consistency» в CAP = линеаризуемость, очень сильная модель. Это НЕ «C» в ACID. Многие «AP»-системы всё равно дают сильноватые гарантии на ключ; многие «CP»-системы настраиваемы. Избегай бинарных ярлыков; говори о ручках настройки.

PACELC — the more useful version

PACELC — более полезная версия

CAP ignores the cost you pay when the system is healthy. PACELC fixes that:

CAP игнорирует цену, которую ты платишь, когда система здорова. PACELC это исправляет:

if (Partition): trade A vs C else (Else): trade L vs C (Latency vs Consistency)
if (Partition): trade A vs C else (Else): trade L vs C (Задержка vs Консистентность)
SystemOn partition (PA/PC)Else, normal (EL/EC)Class
Dynamo / CassandraPA (stay up)EL (fast, eventual)PA/EL
SpannerPC (consistent)EC (TrueTime waits)PC/EC
etcd / ZooKeeperPCECPC/EC
MongoDB (default)PCEC (primary reads)PC/EC
Cassandra (tuned QUORUM)PC-ishEC-ishtunable
СистемаПри разделении (PA/PC)Иначе, норма (EL/EC)Класс
Dynamo / CassandraPA (остаться в сети)EL (быстро, eventual)PA/EL
SpannerPC (консистентно)EC (TrueTime ждёт)PC/EC
etcd / ZooKeeperPCECPC/EC
MongoDB (default)PCEC (чтение с primary)PC/EC
Cassandra (tuned QUORUM)PC-ishEC-ishнастраиваема
Why PACELC matters for DE 99% of the time there is no partition, so the EL vs EC trade is what you actually live with daily: synchronous quorum reads cost p99 latency; eventual reads are fast but may be stale. For an analytics serving layer you usually pick EL and design downstream to tolerate staleness.
Почему PACELC важен для DE 99% времени разделения нет, поэтому компромисс EL vs EC — это то, с чем ты живёшь каждый день: синхронное чтение с кворума стоит p99 задержки; eventual чтения быстрые, но могут быть устаревшими. Для аналитического serving-слоя обычно выбираешь EL и проектируешь downstream, чтобы терпеть staleness.
Interviewer will probe "Spanner claims CA / 'beats CAP'." — It does not. It is CP: under a true partition the minority becomes unavailable. Its trick is making partitions rare (Google's private network) and using TrueTime (bounded clock uncertainty) to keep latency acceptable while staying externally consistent.
О чём спросит интервьюер «Spanner заявляет CA / "побеждает CAP".» — Нет. Это CP: при настоящем разделении меньшинство становится недоступным. Его трюк — делать разделения редкими (частная сеть Google) и использовать TrueTime (ограниченная неопределённость часов), чтобы держать задержку приемлемой, оставаясь externally consistent.

02 Consistency modelsМодели консистентности

A spectrum from "always correct, slow" to "fast, surprising". Know where each sits and the client guarantees.

Спектр от «всегда корректно, медленно» до «быстро, сюрпризы». Знай, где находится каждая и какие гарантии клиенту.

STRONGER / SLOWER WEAKER / FASTER Linearizable ─ Sequential ─ Causal ─ Read-your-writes ─ Eventual (single, (global (cause (session (converges real-time order, no before sees own eventually, order) real-time) effect) writes) no order)
СИЛЬНЕЕ / МЕДЛЕННЕЕ СЛАБЕЕ / БЫСТРЕЕ Linearizable ─ Sequential ─ Causal ─ Read-your-writes ─ Eventual (единый, (глобальный (причина (сессия (сходится real-time порядок, раньше видит свои в итоге, порядок) не real-time) эффекта) записи) нет порядка)
ModelGuaranteeCostUse it for
Linearizable / StrongEvery read returns the latest committed write as of a real-time point; behaves like one copy.Quorum / consensus round-trips; latency + availability hit.Locks, leader election, counters, balances, config.
SequentialAll nodes see ops in same order, but not necessarily real-time order.Cheaper than linearizable.Replicated state machines.
CausalOperations that are causally related are seen in order by everyone; concurrent ops can differ.Track causality (version vectors); cheap, available.Collaborative apps, comments/replies, messaging.
Read-your-writes (session)A client always sees its own prior writes.Sticky routing or version tokens."I posted but don't see it" UX bugs.
Monotonic readsTime never goes backwards for a client.Pin to a replica / hold a token.Avoid "now you see it, now you don't".
EventualIf writes stop, all replicas converge.Cheapest, most available.Caches, analytics, feeds, metrics.
МодельГарантияЦенаПрименяй для
Linearizable / StrongКаждое чтение возвращает последнюю закоммиченную запись на момент real-time; ведёт себя как одна копия.Кворум / консенсус round-trips; удар по задержке + доступности.Блокировки, выбор лидера, счётчики, балансы, конфиги.
SequentialВсе ноды видят операции в одинаковом порядке, но не обязательно real-time.Дешевле linearizable.Реплицированные state machines.
CausalПричинно связанные операции видны всеми в порядке; конкурентные операции могут отличаться.Отслеживание причинности (version vectors); дёшево, доступно.Коллаборативные приложения, комментарии/ответы, мессенджинг.
Read-your-writes (сессия)Клиент всегда видит свои предыдущие записи.Sticky routing или version tokens.Баги UX «я запостил, но не вижу».
Monotonic readsВремя никогда не идёт назад для клиента.Привязка к реплике / держать токен.Избегать «то вижу, то не вижу».
EventualЕсли записи остановятся, все реплики сойдутся.Дешевле всего, максимально доступно.Кеши, аналитика, ленты, метрики.
Gotcha — "eventual" is underspecified Eventual consistency gives no session guarantees by default. A user can write, then read an older replica and not see it. Read-your-writes, monotonic reads, and causal are the practical "session guarantees" that make eventual systems usable for humans.
Ошибка — «eventual» недоспецифицирована Eventual consistency не даёт гарантий сессии по умолчанию. Пользователь может записать, потом прочитать старую реплику и не увидеть. Read-your-writes, monotonic reads и causal — это практические «гарантии сессии», делающие eventual-системы юзабельными для людей.
Interviewer will probe "A user updates their profile photo and on refresh sees the old one — what's wrong and how do you fix it?" → Eventual store + read hit a lagging replica. Fix: read-your-writes via sticky session, write-through cache, or a version token that forces reading a replica that is ≥ the write.
О чём спросит интервьюер «Пользователь обновил фото профиля и при обновлении видит старое — что не так и как чинить?» → Eventual хранилище + чтение попало в отстающую реплику. Лечение: read-your-writes через sticky session, write-through cache или version token, который заставляет прочитать реплику ≥ записи.

03 Partitioning / sharding strategiesСтратегии партиционирования / шардинга

How you split data across nodes determines hotspots, rebalancing pain, and range-scan ability.

Как ты делишь данные по нодам, определяет hotspots, боль ребалансировки и возможность range-сканов.

Hash partitioning

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

node = hash(key) % N. Even spread, kills hotspots, but destroys range queries and re-mods everything when N changes.

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

even loadровная нагрузка no range scansнет range-сканов rehash on resizeперехеш при resize

Range partitioning

Range-партиционирование

Contiguous key ranges per shard (e.g. by date/userid range). Great for scans, but sequential keys create hotspots (the "monotonic timestamp" trap).

Смежные диапазоны ключей на шард (например, по дате/userid range). Отлично для сканов, но последовательные ключи создают hotspots (ловушка «монотонного timestamp»).

range scansrange-сканы needs splitsнужны splits hotspot on hot rangehotspot на горячем range

Consistent hashing (the senior answer to "how do you add a node?")

Consistent hashing (senior-ответ на «как добавить ноду?»)

Keys and nodes both hash onto a ring; a key is owned by the next node clockwise. Adding/removing a node only moves ~K/N keys, not all of them. Virtual nodes (many ring tokens per physical node) smooth out skew and make rebalancing proportional.

Ключи и ноды хешируются на кольцо; ключ принадлежит следующей ноде по часовой стрелке. Добавление/удаление ноды перемещает только ~K/N ключей, а не все. Virtual nodes (много ring-токенов на физическую ноду) сглаживают перекос и делают ребалансировку пропорциональной.

[N1] / \ key hashes here ─┐ [N4] [N2] ▼ \ / owner = next node clockwise = N2 [N3] add N5 between N1,N2 → only keys in that arc move (N2's share), everyone else untouched.
[N1] / \ ключ хешируется сюда ─┐ [N4] [N2] ▼ \ / владелец = следующая нода по часовой = N2 [N3] добавить N5 между N1,N2 → только ключи в этой дуге переезжают (доля N2), остальные не тронуты.
Why vnodes Without vnodes, removing a node dumps ALL its load onto a single neighbor. With ~128–256 vnodes/host, that node's keys spread across the whole cluster, and heterogeneous hardware can carry proportionally more tokens.
Зачем vnodes Без vnodes удаление ноды сваливает ВСЮ её нагрузку на одного соседа. С ~128–256 vnodes/хост ключи этой ноды распределяются по всему кластеру, и гетерогенное железо может нести пропорционально больше токенов.
Interviewer will probe "Your sharded store is hot on one shard." → Diagnose: skewed key (celebrity/whale), low-cardinality partition key, or sequential key. Fixes: salt the key (suffix bucket), composite key, switch to consistent hashing + vnodes, or split the hot range. Tie it to Spark: same fix as a skewed join key (salting).
О чём спросит интервьюер «Твой шардированный store горячий на одном шарде.» → Диагноз: перекошенный ключ (celebrity/whale), низкокардинальный ключ партиционирования или последовательный ключ. Лечение: посолить ключ (suffix bucket), составной ключ, переключиться на consistent hashing + vnodes, или разбить горячий range. Связать со Spark: тот же фикс, что для перекошенного join-ключа (salting).

04 Replication & quorumsРепликация и кворумы

Copies for durability + availability. The interesting part is what happens to writes during failures.

Копии ради надёжности + доступности. Интересная часть — что происходит с записями во время сбоев.

Topologies: leader-follower, multi-leader, leaderless Топологии: leader-follower, multi-leader, leaderless
TopologyHowProsCons / failure mode
Leader–follower (single-leader)All writes to one leader, replicated to followers (sync or async).Simple, no write conflicts, easy strong reads from leader.Leader is bottleneck + SPOF; failover window; async followers lag (stale reads); split-brain risk.
Multi-leaderMultiple leaders accept writes (e.g. per-region), replicate to each other.Low-latency local writes, survives region partition.Write conflicts need resolution (LWW, CRDTs, app merge); convergence complexity.
Leaderless (Dynamo-style)Client writes to many replicas; reads from many; quorum decides.No failover, highly available, tunable.Read repair / anti-entropy needed; you reason about quorums yourself.
ТопологияКакПлюсыМинусы / режим сбоя
Leader–follower (single-leader)Все записи в одного лидера, реплицируются фолловерам (sync или async).Просто, нет конфликтов записи, легко сильные чтения с лидера.Лидер — узкое место + SPOF; окно failover; async-фолловеры отстают (stale reads); риск split-brain.
Multi-leaderНесколько лидеров принимают записи (например, по региону), реплицируют друг другу.Низкая задержка локальных записей, переживает разделение региона.Конфликты записи нужно разрешать (LWW, CRDTs, app merge); сложность сходимости.
Leaderless (Dynamo-стиль)Клиент пишет во много реплик; читает из многих; кворум решает.Нет failover, высокая доступность, настраиваемо.Нужны read repair / anti-entropy; рассуждаешь о кворумах сам.
Sync vs async replication Sync = leader waits for follower ack → no data loss on leader crash, but latency + a slow follower stalls writes. Async = fast, but a leader crash loses un-replicated writes (data loss window). Semi-sync (wait for ≥1 follower) is the common compromise.
Sync vs async репликация Sync = лидер ждёт ack фолловера → нет потери данных при падении лидера, но задержка + медленный фолловер блокирует записи. Async = быстро, но падение лидера теряет нереплицированные записи (окно потери данных). Semi-sync (ждать ≥1 фолловера) — частый компромисс.

Quorums: the R + W > N rule

Кворумы: правило R + W > N

With N replicas, read from R, write to W. If R + W > N, the read and write sets overlap by at least one replica, so a read is guaranteed to see the latest acknowledged write (strong-ish consistency).

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

N = 3 W=2, R=2 → R+W=4 > 3 ✓ overlap guaranteed (balanced) W=3, R=1 → fast reads, slow/fragile writes (all must ack) W=1, R=3 → fast writes, slow reads W=1, R=1 → R+W=2 ≤ 3 ✗ may read stale (pure eventual)
N = 3 W=2, R=2 → R+W=4 > 3 ✓ пересечение гарантировано (баланс) W=3, R=1 → быстрые чтения, медленные/хрупкие записи (все должны ответить) W=1, R=3 → быстрые записи, медленные чтения W=1, R=1 → R+W=2 ≤ 3 ✗ можно прочитать stale (чистый eventual)
Quorum ≠ linearizable R+W>N gives overlap but NOT full linearizability with concurrent writes, sloppy quorums, or read-repair races. True linearizability needs consensus (Raft/Paxos) or a single leader with synchronous commit.
Кворум ≠ линеаризуемый R+W>N даёт пересечение, но НЕ полную линеаризуемость при конкурентных записях, sloppy quorums или гонках read-repair. Настоящая линеаризуемость требует консенсус (Raft/Paxos) или single leader с синхронным commit.
Interviewer will probe "N=3, set W=2,R=2. One replica is down — can you still serve?" → Yes: writes hit the 2 up replicas (W=2 met), reads hit 2 (R=2 met). "Now two are down?" → Writes fail (can't reach W=2), reads fail. Then discuss sloppy quorum + hinted handoff to stay available at the cost of temporary inconsistency.
О чём спросит интервьюер «N=3, установить W=2,R=2. Одна реплика упала — можешь ли обслуживать?» → Да: записи попадают в 2 живые реплики (W=2 выполнено), чтения попадают в 2 (R=2 выполнено). «А если две упали?» → Записи не проходят (не достичь W=2), чтения не проходят. Затем обсудить sloppy quorum + hinted handoff, чтобы остаться доступным ценой временной inconsistency.

05 Consensus & leader election (high level)Консенсус и выбор лидера (high level)

You won't implement Raft, but you must explain why consensus is hard and what it buys.

Ты не будешь имплементировать Raft, но должен объяснить, почему консенсус сложен и что он даёт.

The problem & FLP

Проблема и FLP

Consensus = get a set of nodes to agree on one value (e.g. the order of a log entry, or who is leader) despite crashes and message delays. FLP impossibility: in a fully asynchronous network with even one faulty node, no deterministic protocol can guarantee agreement in bounded time. Real systems sidestep this with partial synchrony + timeouts/randomization — they sacrifice liveness (might stall during chaos) but never safety (never disagree).

Консенсус = заставить набор нод согласиться на одно значение (например, порядок записи в лог или кто лидер), несмотря на падения и задержки сообщений. FLP impossibility: в полностью асинхронной сети даже с одной сбоящей нодой никакой детерминированный протокол не может гарантировать согласие за ограниченное время. Реальные системы обходят это через partial synchrony + timeouts/randomization — жертвуют liveness (могут застопориться во время хаоса), но никогда safety (никогда не разойдутся).

Raft (understandable)

Raft (понятный)

  • Roles: follower → candidate → leader.
  • Election: randomized timeout; candidate requests votes; majority wins for a term.
  • Log replication: leader appends; entry commits once a majority replicates it.
  • Safety: only a node with an up-to-date log can win; committed entries never lost.
  • Роли: follower → candidate → leader.
  • Выбор: рандомизированный таймаут; candidate запрашивает голоса; большинство побеждает на term.
  • Репликация лога: лидер добавляет; запись коммитится, когда большинство её реплицирует.
  • Безопасность: только нода с актуальным логом может победить; закоммиченные записи никогда не теряются.

Paxos / Multi-Paxos

Paxos / Multi-Paxos

  • Phases: prepare (promise) → accept; proposers, acceptors, learners.
  • Multi-Paxos elects a stable leader to skip phase-1 per entry (≈ Raft in practice).
  • Provably correct, famously hard to reason about → Raft was designed as the teachable equivalent.
  • Фазы: prepare (promise) → accept; proposers, acceptors, learners.
  • Multi-Paxos выбирает стабильного лидера, чтобы пропускать phase-1 на запись (≈ Raft на практике).
  • Доказуемо корректен, знаменито сложен для понимания → Raft был создан как понятный эквивалент.
Majority quorum Consensus needs a strict majority (e.g. 2-of-3, 3-of-5). That's why these clusters are odd-sized: a partition leaves at most one side with a majority → only that side makes progress → no split-brain. 5 nodes tolerate 2 failures; 3 tolerate 1.
Кворум большинства Консенсус требует строгое большинство (например, 2 из 3, 3 из 5). Поэтому эти кластеры нечётного размера: разделение оставляет максимум одну сторону с большинством → только эта сторона делает progress → нет split-brain. 5 нод переживают 2 сбоя; 3 переживают 1.
Leader election ≠ consensus done Electing a leader is just one use of consensus. The deeper guarantee is fencing: a deposed-but-revived old leader must be stopped. Use fencing tokens (monotonic IDs) so storage rejects writes from a stale leader. This is the #1 source of split-brain corruption.
Выбор лидера ≠ консенсус сделан Выбор лидера — это лишь одно применение консенсуса. Более глубокая гарантия — fencing: свергнутого-но-ожившего старого лидера надо остановить. Используй fencing tokens (монотонные ID), чтобы хранилище отклоняло записи от stale-лидера. Это источник №1 порчи данных из-за split-brain.
Interviewer will probe "Why odd number of nodes?" / "Two nodes both think they're leader — what now?" → Majority can't exist on both sides of a partition; the minority leader can't get quorum to commit, and fencing tokens make storage reject its writes. Old leader's in-flight writes are safely orphaned.
О чём спросит интервьюер «Почему нечётное число нод?» / «Две ноды обе думают, что они лидеры — что теперь?» → Большинство не может существовать на обеих сторонах разделения; лидер меньшинства не может получить кворум для commit, а fencing tokens заставляют хранилище отклонить его записи. In-flight записи старого лидера безопасно осиротели.

06 Fault tolerance & failure modesОтказоустойчивость и режимы сбоев

Catalogue the failures and your standard mitigations.

Каталогизируй сбои и твои стандартные меры.

FailureSymptomDetectionMitigation
Node crash (crash-stop)Process goneMissed heartbeats / lease expiryReplication + failover; replay from durable log
Gray failure / slow nodeUp but degraded — high latency, partial errorsp99 latency, error-rate, hedged requestsHedging, timeouts, eject from pool, circuit breaker
Network partitionSubsets can't talkQuorum loss, heartbeat gapsCP: minority stops; AP: degrade + reconcile; fencing tokens
Clock skewWrong ordering / TTLsNTP drift monitoringLogical clocks (Lamport/vector), bounded uncertainty (TrueTime); never trust wall-clock for ordering
Poison messageOne record kills the consumer repeatedlyRepeated retries on same offsetDead-letter queue, retry cap, skip-and-log
Thundering herd / retry stormFailure → everyone retries at once → meltdownSpike on recoveryExponential backoff + jitter, retry budgets, circuit breaker
Cascading failureOne slow dep stalls callers, then theirs…Saturation propagationBulkheads, timeouts everywhere, backpressure, load shedding
Data corruption / bit rotSilent wrong bytesChecksums, scrubbingPer-block checksums, erasure coding, re-replicate
СбойСимптомОбнаружениеМеры
Падение ноды (crash-stop)Процесс исчезПропущенные heartbeats / истечение leaseРепликация + failover; переигрывание из надёжного лога
Gray failure / медленная нодаРаботает, но деградирует — высокая задержка, частичные ошибкиp99 latency, error-rate, hedged requestsHedging, timeouts, вышвырнуть из pool, circuit breaker
Разделение сетиПодмножества не могут говоритьПотеря кворума, разрывы heartbeatCP: меньшинство останавливается; AP: деградировать + reconcile; fencing tokens
Перекос часовНеправильный порядок / TTLNTP drift monitoringЛогические часы (Lamport/vector), ограниченная неопределённость (TrueTime); никогда не доверять wall-clock для упорядочивания
Poison messageОдна запись убивает консьюмера циклическиПовторные ретраи на одном offsetDead-letter queue, ограничение ретраев, пропустить-и-залогировать
Thundering herd / шторм ретраевСбой → все одновременно ретраятся → коллапсВсплеск при восстановленииExponential backoff + jitter, retry budgets, circuit breaker
Каскадный сбойОдна медленная зависимость блокирует вызывающих, потом их…Распространение saturationBulkheads, timeouts везде, backpressure, load shedding
Порча данных / bit rotТихо неправильные байтыChecksums, scrubbingChecksums на блок, erasure coding, re-replicate
Gray failure is the killer Crash failures are easy — the node is gone, failover fires. Gray failures (a host that's 30% slow, dropping 2% of requests, GC-pausing) evade binary health checks, poison your tail latency, and are the cause of most real incidents. Detect with differential observability (this replica vs peers) and hedged requests (fire a duplicate to another replica after p95, take the first to answer).
Gray failure — убийца Падения легко — нода исчезла, срабатывает failover. Gray failures (хост 30% медленный, роняет 2% запросов, GC-паузит) обходят бинарные health checks, отравляют tail latency и являются причиной большинства реальных инцидентов. Обнаруживать через дифференциальную observability (эта реплика vs peers) и hedged requests (послать дубликат на другую реплику после p95, взять первого ответившего).

Failure model spectrum

Спектр моделей сбоев

crash-stop ⊂ crash-recovery ⊂ omission ⊂ Byzantine (dies, stays (dies, comes (drops/ (lies / arbitrary — dead) back, must delays msgs) needs BFT, ~3f+1) recover state)
crash-stop ⊂ crash-recovery ⊂ omission ⊂ Byzantine (умирает, (умирает, (теряет/ (лжёт / произвольно — остаётся мёртв) возвращается, задерживает нужен BFT, ~3f+1) должна восстановить msg)
Interviewer will probe "Walk me through what happens when a Kafka broker hosting leader partitions dies." → Controller detects via ZK/KRaft session loss → elects new partition leaders from the ISR (in-sync replicas) → producers with acks=all + min.insync.replicas=2 lose nothing; acks=1 may lose the un-replicated tail. Consumers continue from committed offsets. Name the data-loss window explicitly.
О чём спросит интервьюер «Расскажи, что происходит, когда Kafka-брокер с лидер-партициями умирает.» → Controller обнаруживает через потерю ZK/KRaft-сессии → выбирает новых лидеров партиций из ISR (in-sync replicas) → продюсеры с acks=all + min.insync.replicas=2 ничего не теряют; acks=1 могут потерять нереплицированный хвост. Консьюмеры продолжают с закоммиченных offsets. Назови окно потери данных явно.

07 Idempotency & delivery semanticsИдемпотентность и семантики доставки

"Exactly-once" is the most misunderstood phrase in data engineering. Nail this.

«Exactly-once» — самая неправильно понятая фраза в data engineering. Прибей её.

SemanticWhat it meansHowRisk
At-most-once0 or 1 deliveryFire-and-forget, no retryData loss
At-least-once1 or more (dupes possible)Retry until ackedDuplicates → double counting
Exactly-onceEffect applied onceAt-least-once delivery + dedup / idempotent writes / transactionsCost, complexity, scope limits
СемантикаЧто значитКакРиск
At-most-once0 или 1 доставкаFire-and-forget, без ретраевПотеря данных
At-least-once1 или больше (возможны дубли)Ретраить, пока не ackДубликаты → двойной подсчёт
Exactly-onceЭффект применён один разAt-least-once delivery + dedup / идемпотентные записи / транзакцииЦена, сложность, ограничения области
The truth about exactly-once Exactly-once delivery over a network is impossible (the Two Generals problem — you can't know your ack arrived). What's achievable is exactly-once processing / effect: you accept duplicate delivery and make the effect idempotent. The formula:
at-least-once delivery + idempotent writes (upsert by key) } effectively + dedup (idempotency key / offset tracking) } exactly-once + atomic commit (write + offset in one txn) } EFFECT
Правда про exactly-once Exactly-once доставка по сети невозможна (проблема двух генералов — ты не можешь знать, что твой ack дошёл). Достижимо — exactly-once обработка / эффект: ты принимаешь дублирующую доставку и делаешь эффект идемпотентным. Формула:
at-least-once delivery + идемпотентные записи (upsert по ключу) } фактически + dedup (idемпотентность-ключ / offset tracking) } exactly-once + атомарный commit (запись + offset в одной txn) } ЭФФЕКТ

How real systems do it

Как это делают реальные системы

-- Idempotent sink: dedup billions of at-least-once events by event_id
MERGE INTO fct_events AS target
USING staged_batch AS source
  ON target.event_id = source.event_id
WHEN NOT MATCHED THEN INSERT *;
-- replaying the same batch is a no-op → effective exactly-once
-- Идемпотентный sink: дедуп миллиардов at-least-once событий по event_id
MERGE INTO fct_events AS target
USING staged_batch AS source
  ON target.event_id = source.event_id
WHEN NOT MATCHED THEN INSERT *;
-- переигрывание того же батча — no-op → фактически exactly-once
Idempotency design rules Prefer writes that are naturally idempotent: upsert by key, SET x=v (not x = x + 1), PUT to a deterministic path. If you must increment/aggregate, dedupe first by a unique event id, or make the aggregate a deterministic function of the input set so reprocessing is safe.
Правила дизайна идемпотентности Предпочитай записи, которые естественно идемпотентны: upsert по ключу, SET x=v (не x = x + 1), PUT по детерминированному пути. Если надо инкрементить/агрегировать, сначала дедуп по уникальному event id, или сделай агрегат детерминированной функцией входного множества, чтобы переобработка была безопасной.
Interviewer will probe "Your pipeline can deliver the same event twice. How do you avoid double counting in a daily metric over 3.6B events?" → At-least-once is fine; dedup on a stable event_id via MERGE/window-dedup, partition by day so a reprocessed day is fully idempotent (overwrite-partition pattern), and never use non-idempotent accumulators. Mention watermark for late dupes.
О чём спросит интервьюер «Твой пайплайн может доставить одно событие дважды. Как избежать двойного подсчёта в дневной метрике по 3.6B событиям?» → At-least-once нормально; дедуп по стабильному event_id через MERGE/window-dedup, партиционировать по дню, чтобы переобработанный день был полностью идемпотентен (overwrite-partition pattern), и никогда не использовать неидемпотентные аккумуляторы. Упомяни watermark для поздних дублей.

08 Backpressure & flow controlBackpressure и управление потоком

What a healthy system does when the producer outruns the consumer.

Что здоровая система делает, когда продюсер обгоняет консьюмера.

Backpressure = a feedback signal that slows the producer when the consumer can't keep up, instead of letting unbounded buffers grow until OOM. Without it, a transient slowdown becomes a crash.

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

FAST PRODUCER ──▶ [ bounded queue ] ──▶ SLOW CONSUMER │ full? ▼ signal upstream: slow down / block / shed No backpressure → queue grows → memory blows up → cascade With backpressure → producer throttles → graceful degradation
БЫСТРЫЙ ПРОДЮСЕР ──▶ [ ограниченная очередь ] ──▶ МЕДЛЕННЫЙ КОНСЬЮМЕР │ заполнена? ▼ сигнал наверх: замедлить / блокировать / сбросить Нет backpressure → очередь растёт → память взрывается → каскад С backpressure → продюсер тормозит → graceful degradation

Mechanisms (from gentlest to harshest)

Механизмы (от мягчайшего к жёсткому)

  1. Pull-based / demand (Reactive Streams, Flink credit-based): consumer requests N items; producer can't exceed demand. Cleanest.
  2. Bounded buffers + blocking: producer blocks when queue full (TCP flow control works exactly like this).
  3. Rate limiting / token bucket: cap ingress rate at the edge.
  4. Load shedding: drop low-priority work when overloaded (better to drop 1% than fall over).
  5. Buffer to durable log: Kafka as a shock absorber — producer writes at full speed, consumer drains at its own pace; the log is the backpressure-free buffer (bounded by retention/disk).
  1. Pull-based / demand (Reactive Streams, Flink credit-based): консьюмер запрашивает N элементов; продюсер не может превысить запрос. Самый чистый.
  2. Ограниченные буферы + блокировка: продюсер блокируется, когда очередь заполнена (TCP flow control работает именно так).
  3. Rate limiting / token bucket: ограничить ingress rate на краю.
  4. Load shedding: сбросить низкоприоритетную работу при перегрузке (лучше сбросить 1%, чем упасть).
  5. Буферизация в надёжный лог: Kafka как амортизатор — продюсер пишет на полной скорости, консьюмер сливает в своём темпе; лог и есть буфер без backpressure (ограничен retention/disk).
Kafka decouples but doesn't remove the problem A log buffers spikes, but if the consumer is chronically slower than the producer, lag grows unboundedly until retention drops un-consumed data (silent loss). Monitor consumer lag; scale partitions/consumers; the log buys time, not infinite capacity.
Kafka развязывает, но не убирает проблему Лог буферизует всплески, но если консьюмер хронически медленнее продюсера, lag растёт неограниченно, пока retention не сбросит непотреблённые данные (тихая потеря). Мониторь lag консьюмера; масштабируй партиции/консьюмеров; лог покупает время, а не бесконечную ёмкость.
Interviewer will probe "Your streaming job's input rate spikes 5x. What happens?" → With backpressure: Flink/Spark throttle source reads, latency rises but it stays up; lag accumulates in Kafka and drains after the spike. Without: unbounded state/heap → GC death → restart loop. Mention checkpoint duration growing as an early warning.
О чём спросит интервьюер «Input rate твоей streaming-джобы скакнул в 5x. Что происходит?» → С backpressure: Flink/Spark тормозят чтение из source, latency растёт, но она держится; lag копится в Kafka и стекает после всплеска. Без: неограниченное state/heap → GC death → restart loop. Упомяни растущую checkpoint duration как раннее предупреждение.

09 Data skew, hotspots & shufflesПерекос данных, hotspots и shuffle

The bread-and-butter of big-data DE. Your Spark experience shines here.

Хлеб с маслом big-data DE. Здесь сияет твой Spark-опыт.

The shuffle

Shuffle

A shuffle = redistributing data across the network so all rows with the same key land on the same task (needed for joins, groupBy, distinct, repartition). It's the most expensive primitive: disk spill + network + serialization. Wide transformations (join, groupByKey, reduceByKey, distinct) shuffle; narrow ones (map, filter) don't.

Shuffle = перераспределение данных по сети, чтобы все строки с одним ключом попали в одну задачу (нужно для joins, groupBy, distinct, repartition). Это самая дорогая операция: disk spill + network + serialization. Wide-трансформации (join, groupByKey, reduceByKey, distinct) делают shuffle; narrow (map, filter) — нет.

stage 1 (map) SHUFFLE stage 2 (reduce) [t0][t1][t2] ──┐ by hash(key) ┌──▶ [r0] all key-A rows [t0][t1][t2] ──┼───────────────┼──▶ [r1] all key-B rows [t0][t1][t2] ──┘ └──▶ [r2] all key-C rows network + disk spill = the bottleneck
stage 1 (map) SHUFFLE stage 2 (reduce) [t0][t1][t2] ──┐ по hash(key) ┌──▶ [r0] все строки key-A [t0][t1][t2] ──┼────────────────┼──▶ [r1] все строки key-B [t0][t1][t2] ──┘ └──▶ [r2] все строки key-C network + disk spill = узкое место

Data skew = one partition gets most of the rows

Перекос данных = одна партиция получает большинство строк

If one key is 1000x more common (a "whale" / celebrity / null key), one reduce task does most of the work while others idle → the job is gated by that one straggler. This is the same root cause as a sharded-DB hotspot.

Если один ключ в 1000x чаще («whale» / celebrity / null key), одна reduce-задача делает большую часть работы, пока остальные простаивают → джоба ограничена этим одним отстающим. Та же корневая причина, что hotspot в шардированной БД.

FixHowWhen
SaltingAppend random bucket to hot key (key_0..key_N), aggregate in two passes.Skewed groupBy / heavy key
Broadcast joinShip the small side to every executor → no shuffle at all.One side fits in memory
AQE skew joinSpark Adaptive Query Execution auto-splits skewed partitions at runtime.Spark 3+, enable AQE
Isolate / separate the whaleHandle the hot key path separately, union results.Known celebrity keys
Filter nulls / sentinel keysNull join keys all hash to one partition — handle separately.Dirty data
Tune partitionsRight-size shuffle.partitions; avoid tiny-files / mega-partitions.Always
ФиксКакКогда
SaltingДобавить случайный bucket к горячему ключу (key_0..key_N), агрегировать в два прохода.Перекошенный groupBy / тяжёлый ключ
Broadcast joinДоставить малую сторону на каждый executor → shuffle вообще нет.Одна сторона влезает в память
AQE skew joinSpark Adaptive Query Execution автоматически разбивает перекошенные партиции в runtime.Spark 3+, включить AQE
Изолировать / отделить whaleОбработать путь горячего ключа отдельно, union результаты.Известные celebrity-ключи
Фильтровать null / sentinel keysNull join-ключи все хешируются в одну партицию — обработать отдельно.Грязные данные
Тюнить партицииПодобрать размер shuffle.partitions; избегать tiny-files / мега-партиций.Всегда
# Salting a skewed key in Spark to spread one hot key across N tasks
N = 50
salted = df.withColumn("salt", (rand()*N).cast("int")) \
           .groupBy("key", "salt").agg(_sum("v").alias("partial"))
result = salted.groupBy("key").agg(_sum("partial"))  # 2-phase
# Salting перекошенного ключа в Spark для распределения одного горячего ключа на N задач
N = 50
salted = df.withColumn("salt", (rand()*N).cast("int")) \
           .groupBy("key", "salt").agg(_sum("v").alias("partial"))
result = salted.groupBy("key").agg(_sum("partial"))  # 2-phase
Interviewer will probe "A Spark stage has 199 tasks done in 30s and 1 task running for 40 min." → Textbook skew. Diagnose via Spark UI (task input-size distribution). Fix per table above. Connect it: "skew in Spark = hotspot in a sharded store = same salting/splitting cure." This is exactly the senior cross-domain reasoning they want.
О чём спросит интервьюер «У Spark stage 199 задач выполнились за 30s, и 1 задача крутится 40 минут.» → Учебниковый skew. Диагноз через Spark UI (распределение input-size задач). Фиксить по таблице выше. Связать: «skew в Spark = hotspot в шардированном store = то же лечение salting/splitting.» Это именно то кросс-доменное рассуждение senior-уровня, которое они хотят.

10 Designing pipelines for billions of events/dayПроектирование пайплайнов на миллиарды событий/день

Anchor to your real number: ~3.6B events/day ≈ ~42k events/sec average, multiples of that at peak.

Зацепись за твоё реальное число: ~3.6B событий/день ≈ ~42k событий/сек в среднем, кратно больше на пике.

Reference architecture

Эталонная архитектура

producers ─▶ Kafka (partitioned, replicated log) ─▶ stream proc (Flink/SSS) │ │ └──▶ raw landing (S3/HDFS, immutable) └─▶ realtime sink │ batch (Spark) ─▶ staged ─▶ dims/facts ─▶ serving orchestration: Airflow (deps, retries, SLAs, backfills)
продюсеры ─▶ Kafka (партиционированный, реплицированный лог) ─▶ stream proc (Flink/SSS) │ │ └──▶ raw landing (S3/HDFS, immutable) └─▶ realtime sink │ batch (Spark) ─▶ staged ─▶ dims/facts ─▶ serving оркестрация: Airflow (зависимости, ретраи, SLA, backfills)

Design principles for scale + correctness

Принципы дизайна для масштаба + корректности

Capacity sanity-check (do this live) 3.6B/day ÷ 86,400 ≈ 41.7k events/s avg. Assume 5x peak ≈ 200k/s. At ~1KB/event that's ~200 MB/s peak ingress, ~3.6 TB/day raw. Partition Kafka for ≥ peak/(per-partition throughput); size consumer parallelism to that. Showing you can do the back-of-envelope is a senior signal.
Sanity-check по мощности (делай вживую) 3.6B/день ÷ 86,400 ≈ 41.7k событий/с средн.. Предположим пик 5x ≈ 200k/s. При ~1KB/событие это ~200 MB/s пиковый ingress, ~3.6 TB/день raw. Партиционировать Kafka на ≥ пик/(пропускная способность партиции); подобрать параллелизм консьюмеров под это. Показать, что можешь посчитать на салфетке — сигнал senior-уровня.
Lambda vs Kappa Lambda = separate batch + speed layers (accurate batch reconciles fast-but-approximate stream) — robust but two codebases to keep in sync. Kappa = one streaming path, reprocess by replaying the log — simpler, needs strong replay + idempotency. Most modern DE leans Kappa with a batch correction layer only where needed.
Lambda vs Kappa Lambda = отдельные batch + speed слои (точный batch сверяет быстрый-но-приблизительный stream) — надёжно, но две кодовые базы держать в синхроне. Kappa = один streaming-путь, переобработка через переигрывание лога — проще, требует сильное replay + идемпотентность. Современный DE склоняется к Kappa с batch-коррекцией только где нужно.
Interviewer will probe "A bug corrupted yesterday's output. How do you fix it without double counting?" → Because raw is immutable + stages are idempotent + partitioned by date: re-run only that date's partition with overwrite semantics (Airflow backfill of one ds). No dupes, no manual cleanup. This is why idempotency + partitioning are non-negotiable at scale.
О чём спросит интервьюер «Баг испортил вчерашний вывод. Как чинить без двойного подсчёта?» → Потому что raw immutable + стадии идемпотентны + партиционированы по дате: перезапустить только партицию этой даты с overwrite semantics (Airflow backfill одного ds). Нет дублей, нет ручной чистки. Поэтому идемпотентность + партиционирование — неоспоримо на масштабе.

11 Batch vs streaming, watermarks, latency vs throughputBatch vs streaming, watermarks, задержка vs пропускная способность

Tradeoffs you choose deliberately, plus the late-data problem.

Компромиссы, которые выбираешь осознанно, плюс проблема поздних данных.

DimensionBatchStreaming
LatencyMinutes–hoursms–seconds
ThroughputVery high (bulk efficiency)High, but per-record overhead
CorrectnessEasy (complete data, reprocess freely)Hard (late/out-of-order, partial windows)
Cost / complexityCheaper, simpler opsAlways-on, stateful, harder
Use forDaily aggregates, training data, backfillsAlerting, fraud, real-time dashboards
ИзмерениеBatchStreaming
ЗадержкаМинуты–часымс–секунды
Пропускная способностьОчень высокая (bulk efficiency)Высокая, но overhead на запись
КорректностьЛегко (полные данные, свободная переобработка)Сложно (поздние/out-of-order, частичные окна)
Цена / сложностьДешевле, проще opsAlways-on, stateful, сложнее
Применять дляДневные агрегаты, training data, backfillsАлертинг, fraud, дашборды в real-time

Event time vs processing time & watermarks

Event time vs processing time и watermarks

Events arrive out of order and late (mobile offline, retries, partitions). You usually want to aggregate by event time (when it happened), not processing time (when you saw it). A watermark is the system's assertion: "I believe I've seen all events with event-time ≤ W." It's the trade between completeness and latency.

События приходят не по порядку и поздно (мобильные offline, ретраи, партиции). Обычно хочется агрегировать по event time (когда произошло), а не processing time (когда увидел). Watermark — это утверждение системы: «Я считаю, что видел все события с event-time ≤ W.» Это компромисс между полнотой и задержкой.

event time ──▶ ... e@10:00 e@10:02 e@09:58(late!) e@10:03 ... watermark = max_event_time − allowedLateness window [10:00,10:05) fires when watermark passes 10:05 • late event before watermark → still included • late event after watermark → dropped OR side-output (DLQ) OR triggers update
event time ──▶ ... e@10:00 e@10:02 e@09:58(late!) e@10:03 ... watermark = max_event_time − allowedLateness окно [10:00,10:05) срабатывает, когда watermark прошёл 10:05 • позднее событие до watermark → ещё включено • позднее событие после watermark → сброшено ИЛИ side-output (DLQ) ИЛИ триггерит update

Throughput vs latency (general)

Пропускная способность vs задержка (общее)

Batching amortizes per-op cost → higher throughput but adds queueing latency. Little's Law: L = λ × W (items in system = arrival rate × time in system). You can't max both; pick the SLO. For analytics: favor throughput. For serving/alerting: favor latency, cap p99.

Батчинг амортизирует цену на операцию → выше пропускная способность, но добавляется queueing задержка. Little's Law: L = λ × W (элементы в системе = частота прибытия × время в системе). Нельзя максимизировать оба; выбирай SLO. Для аналитики: предпочитай пропускную способность. Для serving/alerting: предпочитай задержку, ограничивай p99.

Interviewer will probe "Events can arrive 2 hours late. How do you count daily actives correctly?" → Aggregate by event time, set watermark/allowed-lateness to cover the realistic tail, accept a small drop beyond it, AND run a batch reprocess of the day after the window closes to true-up. State the explicit trade: real-time number is approximate, batch number is final.
О чём спросит интервьюер «События могут приходить на 2 часа позже. Как правильно посчитать daily actives?» → Агрегировать по event time, установить watermark/allowed-lateness, чтобы покрыть реалистичный хвост, принять небольшой drop за его пределами, И запустить batch-переобработку дня после закрытия окна для уточнения. Чётко озвучь компромисс: real-time число приблизительное, batch-число финальное.
Don't use wall-clock for ordering Processing-time windows give wrong results whenever ingestion lags or replays. Always key correctness on a trustworthy event_time in the payload, and treat the watermark as your tunable completeness/latency dial.
Не используй wall-clock для упорядочивания Processing-time окна дают неправильные результаты, когда ingestion отстаёт или переигрывается. Всегда привязывай корректность к надёжному event_time в payload, а watermark трактуй как настраиваемую ручку полноты/задержки.

12 Self-quiz — flip to revealСамопроверка — переверни карту

Tap any card. Answer out loud first, then flip.

Нажми на любую карту. Сначала ответь вслух, потом переверни.

CAP / PACELC
Why is "CA" not a real design point, and what does PACELC add over CAP?Почему «CA» не реальная точка дизайна, и что PACELC добавляет к CAP?
tap to flipнажми
Partitions are inevitable in a networked system, so you can't opt out of P — you only choose C-vs-A during a partition. PACELC adds the Else case: when healthy (the common case) you still trade Latency vs Consistency. That EL/EC trade is what you actually live with daily.Разделения неизбежны в сетевой системе, поэтому не можешь отказаться от P — выбираешь только C-против-A во время разделения. PACELC добавляет случай Else: когда здорова (обычный случай), всё равно торгуешь Задержка vs Консистентность. Этот EL/EC компромисс — то, с чем живёшь каждый день.
Delivery semanticsСемантики доставки
Is exactly-once delivery possible? How do real systems get exactly-once effects?Возможна ли exactly-once доставка? Как реальные системы получают exactly-once эффекты?
tap to flipнажми
Exactly-once delivery over a network is impossible (Two Generals — can't confirm the ack). You achieve exactly-once effect = at-least-once delivery + idempotent writes (upsert by key) + dedup (idempotency key / offset tracking) + atomic commit of write+offset. Kafka EOS = idempotent producer + transactions + read_committed.Exactly-once доставка по сети невозможна (Two Generals — не можешь подтвердить ack). Достигаешь exactly-once эффект = at-least-once delivery + идемпотентные записи (upsert по ключу) + dedup (idempотency key / offset tracking) + atomic commit запись+offset. Kafka EOS = идемпотентный продюсер + транзакции + read_committed.
QuorumsКворумы
N=5. Pick R and W for strong-ish reads that survive 2 node failures on the write path. Why?N=5. Выбери R и W для strong-ish чтений, которые переживают 2 падения нод на write path. Почему?
tap to flipнажми
W=3, R=3 → R+W=6 > 5 (overlap guaranteed) and a majority. W=3 tolerates 2 failures (3 of 5 still reachable). Note quorum overlap ≠ full linearizability with concurrent writes/read-repair races — for that you need consensus or a single sync-committed leader.W=3, R=3 → R+W=6 > 5 (пересечение гарантировано) и большинство. W=3 переживает 2 сбоя (3 из 5 всё ещё доступны). Заметь: пересечение кворума ≠ полная линеаризуемость при конкурентных записях/гонках read-repair — для этого нужен консенсус или single sync-committed лидер.
ConsensusКонсенсус
During a partition, why can't both sides keep a leader, and what stops a revived old leader from corrupting data?При разделении, почему обе стороны не могут держать лидера, и что останавливает ожившего старого лидера от порчи данных?
tap to flipнажми
A leader must hold a majority quorum to commit; only one side of a partition can have a majority, so the minority can't make progress → no split-brain. A revived stale leader is stopped by fencing tokens (monotonic IDs): storage rejects any write carrying an older token. Odd cluster sizes guarantee a unique majority.Лидер должен держать кворум большинства для commit; только одна сторона разделения может иметь большинство, поэтому меньшинство не может делать progress → нет split-brain. Ожившего stale-лидера останавливают fencing tokens (монотонные ID): хранилище отклоняет любую запись со старым токеном. Нечётные размеры кластера гарантируют уникальное большинство.
SkewПерекос
One Spark reduce task runs 50x longer than the rest. Diagnose and give two fixes.Одна Spark reduce-задача крутится в 50x дольше остальных. Диагноз и два фикса.
tap to flipнажми
Data skew: one key (whale/null/sentinel) hashes to one partition. Confirm via Spark UI task input-size distribution. Fixes: (1) salt the hot key + two-phase aggregate; (2) broadcast join if the other side is small; (3) enable AQE skew join; (4) isolate the whale / filter null keys. Same cure as a sharded-DB hotspot.Перекос данных: один ключ (whale/null/sentinel) хешируется в одну партицию. Подтверди через Spark UI распределение input-size задач. Фиксы: (1) salt горячий ключ + two-phase aggregate; (2) broadcast join, если другая сторона малая; (3) включить AQE skew join; (4) изолировать whale / отфильтровать null keys. То же лечение, что hotspot в шардированной БД.
StreamingСтриминг
What is a watermark and what does tuning it cost you?Что такое watermark и какова цена его настройки?
tap to flipнажми
A watermark = the system's claim "I've seen all events with event-time ≤ W"; it gates when event-time windows fire. Tighter watermark → lower latency but more dropped late data. Looser → more complete but higher latency + more retained state. Very-late events go to side-output/DLQ or a batch true-up pass.Watermark = утверждение системы «Я видел все события с event-time ≤ W»; ограничивает, когда срабатывают event-time окна. Жёсткий watermark → ниже задержка, но больше сброшенных поздних данных. Свободный → полнее, но выше задержка + больше держимого state. Очень поздние события идут в side-output/DLQ или batch true-up pass.
Failure modesРежимы сбоев
Why is a gray failure harder to handle than a crash?Почему gray failure сложнее обработать, чем падение?
tap to flipнажми
A crash is binary — health checks fail, failover fires. A gray failure (slow/partially-erroring but "up") passes naive health checks while poisoning tail latency and silently dropping a fraction of work. Detect with differential metrics (this replica vs peers) + hedged requests; mitigate by ejecting from the pool, timeouts, and circuit breakers.Падение бинарно — health checks падают, срабатывает failover. Gray failure (медленная/частично ошибающаяся, но «работает») проходит наивные health checks, отравляя tail latency и тихо роняя часть работы. Обнаруживай через дифференциальные метрики (эта реплика vs peers) + hedged requests; лечи выкидыванием из pool, таймаутами и circuit breakers.
Backpressure
Input rate 5x's. With and without backpressure — what happens, and where's the silent data loss risk?Input rate скакнул в 5x. С и без backpressure — что происходит и где риск тихой потери данных?
tap to flipнажми
With: sources throttle, latency rises, lag buffers in Kafka and drains post-spike — stays up. Without: unbounded buffers/state → OOM → restart loop → cascade. Silent loss: if the consumer is chronically slower, Kafka lag exceeds retention and un-consumed records are deleted. Monitor consumer lag + checkpoint duration.С: sources тормозят, latency растёт, lag буферизуется в Kafka и стекает после всплеска — держится. Без: неограниченные buffers/state → OOM → restart loop → каскад. Тихая потеря: если консьюмер хронически медленнее, Kafka lag превышает retention и непотреблённые записи удаляются. Мониторить consumer lag + checkpoint duration.

13 One-liners to deploy under pressureОднострочники для использования под давлением

Framing for Senior/Lead DE Tie answers to scale + reliability of data pipelines: correctness under partial failure, predictable p99, idempotent reprocessing, and observability of gray failures. Lead with the failure-mode reasoning checklist from section 00.
Фрейминг для Senior/Lead DE Привязывай ответы к масштабу + надёжности пайплайнов данных: корректность при частичном сбое, предсказуемый p99, идемпотентная переобработка и observability gray failures. Начинай с чек-листа рассуждений о режимах сбоев из раздела 00.