% distributed-system-design.tex — the backend toolkit a senior mobile engineer
% is expected to speak in a system-design round: load balancers, cache layers
% and patterns, replication + lag, sharding + hot keys, CAP/PACELC, consistency
% models, queues, back-of-envelope numbers, the canonical architecture.
% Source: docs/memos/system-design-distributed.md (checked against
% docs/school/notes/knowledge-gaps-2026-09-23.md). Latency table: the classic
% "numbers every programmer should know", rounded to orders of magnitude.
% Already printed, only cross-referenced: interview/ios-system-design-deep-dives.tex
% (client side), design/resilience-patterns.tex (timeouts, retries, breakers,
% rate limits), design/event-driven-cqrs-sourcing.tex (outbox, sagas),
% cs/hashing-deep.tex (consistent-hash ring).
% Build ONLY with: tools/print/print-sheet.py <this>.tex --dry-run
% @source: hiot monorepo, docs/school/sheets/cs/distributed-system-design.tex — the SOURCE OF TRUTH; a copy anywhere else (e.g. artur.gurgul.pro) is regenerated from it, never edited
% @labels: area=cs kind=architecture level=senior platform=backend new=no round=missing-2026-09-25 topic=system-design
% @tags: load-balancer, cache-aside, cache-stampede, read-replicas, replication-lag, sharding, hot-key, cap-theorem, pacelc, consistency-models, message-queue, back-of-envelope
\documentclass[8pt]{extarticle}
\usepackage{printup-sheet}
\usepackage{array}

\newcommand\ct[1]{\texttt{#1}}
\tikzset{
  s/.style={box, font=\tiny, minimum height=5.5mm, inner sep=1.5pt},
  g/.style={s, draw=sheetGreen, fill=sheetGreen!10},
  r/.style={s, draw=sheetRed, fill=sheetRed!7},
  o/.style={s, draw=sheetOrange, fill=sheetOrange!10},
  gr/.style={s, draw=sheetGrey, fill=black!5},
  hd/.style={font=\bfseries\small, text=sheetBlue, anchor=west},
  l/.style={font=\tiny, text=black!75, align=center, inner sep=1pt},
  db/.style={cylinder, shape border rotate=90, aspect=0.25, draw=sheetBlue, thick,
             fill=sheetBlue!8, font=\tiny, minimum height=8mm, minimum width=11mm, align=center, inner sep=1pt},
}

\begin{document}

\sheettitle{Distributed system design — the toolkit}{cs · memo}

\oneliner{Scale by making servers \textbf{stateless} and adding them behind a
\textbf{load balancer}; take read load off the database with \textbf{caches}
and \textbf{read replicas}; take write and storage load with
\textbf{sharding}; move slow work onto \textbf{queues}. Each step buys scale
with a named cost — staleness, replication lag, hot keys, duplicates — and a
senior answer names the cost and the \textbf{number} that justifies it.}

\vspace{2pt}
\noindent\begin{tikzpicture}[sheet]
  \node[hd] at (-0.1,2.45) {The canonical read-heavy architecture — and what each hop costs};
  \node[s, minimum width=11mm] (cl) at (0.45,0.9) {iOS app\\HTTP cache};
  \node[gr, minimum width=11mm] (cdn) at (2.0,1.75) {CDN edge\\origin: object store};
  \node[s, minimum width=12mm] (lb) at (2.0,0.25) {L7 load\\balancer (×2)};
  \draw[hot] (cl) -- node[l, above left]{DNS} (cdn);
  \draw[hot] (cl) -- node[l, below left]{TLS} (lb);
  \foreach \y [count=\i] in {1.15,0.25,-0.65} \node[g, minimum width=12mm] (a\i) at (4.0,\y) {app server \i\\stateless};
  \foreach \i in {1,2,3} \draw[flow] (lb.east) -- (a\i.west);
  \node[o, minimum width=12mm] (rc) at (6.3,1.75) {Redis\\cache-aside};
  \node[db] (ld) at (6.3,0.2) {leader\\writes};
  \node[db, fill=sheetBlue!3] (f1) at (8.3,0.75) {follower\\reads};
  \node[db, fill=sheetBlue!3] (f2) at (8.3,-0.45) {follower\\reads};
  \draw[flow] (a1.east) -- (rc.west);
  \draw[flow] (a2.east) -- node[l, above]{write} (ld.west);
  \draw[flow, dashed] (a3.east) to[out=0,in=-150] node[l, below, pos=0.6]{read} (f2.south west);
  \draw[hot, sheetRed] (ld.east) -- node[l, above, sloped]{async} (f1.west);
  \draw[hot, sheetRed] (ld.east) -- (f2.west);
  \node[l, text=sheetRed, align=left, anchor=west] at (9.0,1.55) {replication lag:\\ms … seconds};
  \node[s, minimum width=12mm, draw=sheetBrown, fill=sheetBrown!8] (q) at (4.0,-1.55) {queue\\(at-least-once)};
  \draw[flow] (a3.south) -- (q.north);
  \node[s, minimum width=11mm, draw=sheetBrown, fill=sheetBrown!8] (w) at (6.3,-1.55) {workers\\idempotent};
  \draw[flow] (q) -- (w);
  \node[gr, minimum width=11mm] (s3) at (8.3,-1.55) {object store\\S3 / blobs};
  \draw[flow] (w) -- (s3);
  % side panel: sharding + hot key
  \draw[sheetGrey!50] (10.45,2.6) -- (10.45,-1.95);
  \node[hd] at (10.5,2.1) {Sharding by key};
  \node[s, minimum width=15mm] (rt) at (12.9,1.4) {router: shard = f(user\_id)};
  \foreach \x/\n/\c [count=\i] in {11.2/{S0\\a–f}/sheetBlue, 12.3/{S1\\g–m}/sheetBlue, 13.4/{S2\\n–s}/sheetRed, 14.5/{S3\\t–z}/sheetBlue} {
    \node[db, minimum width=9mm, draw=\c, fill=\c!8] (sh\i) at (\x,0.2) {\n};
    \draw[flow] (rt.south) -- (sh\i.north); }
  \node[l, text=sheetRed, align=left, anchor=west] at (10.5,-0.75) {\textbf{hot key}: one celebrity's\\writes all land on S2};
  \node[l, align=left, anchor=west] at (13.1,-0.75) {fix: split the key\\(\ct{id\#0…\#9}), cache,\\own shard};
  \node[l, align=left, anchor=west] at (10.5,-1.55) {range → hot spots on monotonic keys (time);\\hash → even, but no range scans;\\many fixed logical partitions → cheap resharding};
\end{tikzpicture}

\begin{multicols}{2}
\footnotesize\setstretch{1.0}

\section{Load balancers}
{\scriptsize
\begin{tabular}{@{}>{\raggedright\arraybackslash}p{7mm}>{\raggedright\arraybackslash}p{32mm}>{\raggedright\arraybackslash}p{33mm}@{}}
\toprule
& \textbf{sees} & \textbf{can}\\
\midrule
\textbf{L4} & IP + port, TCP/UDP flows & fast, protocol-agnostic, TLS passthrough; no per-request routing\\
\textbf{L7} & HTTP: path, headers, cookies & terminate TLS, route \ct{/api} vs \ct{/img}, retries, rate limits, sticky sessions\\
\bottomrule
\end{tabular}}\par
Algorithms: round-robin (weighted), \textbf{least connections}, least
response time, IP/consistent hash (affinity, cache locality), \textbf{power
of two choices} (pick 2 at random, send to the less loaded — near-optimal,
no global state). Health checks eject bad nodes. The LB itself is redundant
(active–passive pair + floating IP, or DNS/anycast). Session state in a
server's memory breaks all of it → \textbf{stateless} servers, state in
Redis/DB or a signed token.

\section{Caching — layers and patterns}
Layers: app/HTTP cache (\ct{Cache-Control}, \ct{ETag}) → CDN → reverse proxy →
in-process → distributed (Redis/Memcached) → DB buffer pool.
{\scriptsize
\begin{tabular}{@{}>{\raggedright\arraybackslash}p{15mm}>{\raggedright\arraybackslash}p{33mm}>{\raggedright\arraybackslash}p{24mm}@{}}
\toprule
\textbf{cache-aside} & app: read cache; miss → DB → fill; write DB, then \textbf{delete} key & default; stale window small\\
read-through & cache loads from DB itself & simpler app code\\
write-through & write cache + DB synchronously & fresh; slower writes\\
\textbf{write-back} & write cache, flush to DB later & fastest; \textbf{loses data} on crash\\
write-around & write DB only; cache on read & no churn for write-once data\\
\bottomrule
\end{tabular}}\par
\textbf{Invalidate}: TTL (bound staleness) + delete-on-write (not update: two
racing writers can leave the old value) or versioned keys. \textbf{Stampede}
(hot key expires, 1000 misses hit the DB): single-flight / request
coalescing, a lock + serve stale, early probabilistic refresh, TTL jitter.
Effective latency $= h \cdot t_{cache} + (1-h) \cdot t_{db}$: 90\,\% → 99\,\% hit
rate cuts DB load \textbf{10$\times$}.

\section{Databases — replication and sharding}
\begin{itemize}
  \item \textbf{Leader–follower}: leader takes writes, streams its log.
        \textbf{Async}: fast, but followers \textbf{lag} and a failover can lose the
        last writes; \textbf{sync}/semi-sync: durable, slower. Read replicas scale
        \emph{reads} only.
  \item Lag breaks \textbf{read-your-writes} (post, refresh, it's gone): read
        your own data from the leader for a few seconds, or wait until the
        replica passed your write's log position. \textbf{Monotonic reads}: stick a
        user to one replica.
  \item \textbf{Multi-leader} (multi-region, offline clients): write conflicts →
        last-write-wins (loses data) or CRDTs/merge. \textbf{Leaderless} (Dynamo,
        Cassandra): N replicas, write W, read R; \textbf{R + W > N} ⇒ read and write
        sets overlap.
  \item \textbf{Sharding} splits writes + storage; costs: cross-shard joins and
        transactions (denormalise, sagas), resharding, hot keys (panel above).
        Shard key = the access pattern (\ct{user\_id}, \ct{space\_id}).
  \item Indexes (B-tree) make reads O(log n) and every write slower. \ct{EXPLAIN}.
\end{itemize}

\columnbreak

\section{CAP and PACELC, in plain words}
\textbf{CAP}: when the network \textbf{partitions}, a replica cut off from the
others must either refuse (\textbf{C}onsistent) or answer with possibly stale
data (\textbf{A}vailable). P is not optional, so it is \textbf{CP vs AP during
a partition} — not ``pick 2 of 3''. \textbf{PACELC}: \emph{if P}, A or C;
\emph{else} (normal operation) \textbf{L}atency or \textbf{C}onsistency — the
trade you actually pay every day (sync replication = slower writes).

\section{Consistency models — strongest first}
{\scriptsize
\begin{tabular}{@{}>{\raggedright\arraybackslash}p{19mm}>{\raggedright\arraybackslash}p{53mm}@{}}
\toprule
\textbf{linearizable} & behaves like one copy; a read sees every completed write (bank balance, locks, unique usernames)\\
sequential / causal & everyone sees causally related writes in order (reply after the post)\\
\textbf{read-your-writes} & \emph{you} always see your own writes (profile edit)\\
monotonic reads & time never goes backwards for one user\\
\textbf{eventual} & replicas converge if writes stop (like counts, feeds)\\
\bottomrule
\end{tabular}}

\section{Queues — async work}
Decouple producer from consumer, absorb bursts, fan out. Delivery is
\textbf{at-least-once} → consumers \textbf{idempotent} (dedupe on a message
id). Poison messages → retry with backoff → \textbf{dead-letter queue}.
Ordering only per partition key (Kafka partition). \textbf{Log} (Kafka:
retained, replayable, many readers) vs \textbf{queue} (SQS/RabbitMQ: deleted on
ack). DB write + publish atomically → \textbf{outbox}. Watch \textbf{queue
depth / age} — it is your backlog alarm.

\section{Back-of-envelope — numbers to know}
{\scriptsize
\begin{tabular}{@{}>{\raggedright\arraybackslash}p{34mm}>{\raggedright\arraybackslash}p{19mm}>{\raggedright\arraybackslash}p{18mm}@{}}
\toprule
L1 cache / main memory reference & 1 / 100 ns & RAM $\approx$ 100$\times$ L1\\
SSD random read (4 KB) & $\sim$100\,µs & 1000$\times$ RAM\\
read 1 MB: RAM / SSD & $\sim$10s µs / $\sim$1 ms & disk seek 5–10 ms\\
round trip in one datacenter & 0.5 ms & \\
RTT same continent / cross-Atlantic & 20–50 / $\sim$80 ms & mobile LTE $\sim$50 ms\\
frame budget 60 / 120 Hz & 16.7 / 8.3 ms & \\
\bottomrule
\end{tabular}}\par
\textbf{Method}: 1 day $\approx$ 86\,400 s $\approx$ 10\textsuperscript{5} s → 1M
requests/day $\approx$ 12 QPS. \textbf{QPS} = DAU $\times$ requests per user ÷
10\textsuperscript{5}; \textbf{peak} = 2–3$\times$ average (more for events).
\textbf{Storage} = items/day $\times$ size $\times$ retention $\times$ 3 replicas.
State assumptions, round hard, keep the order of magnitude.

\textbf{Example.} 10M DAU $\times$ 20 requests = 2$\times$10\textsuperscript{8}/day
→ $\sim$2\,300 QPS avg, $\sim$7\,000 peak. 10\,\% post one 500\,KB photo/day →
500 GB/day → $\sim$180 TB/year, $\times$3 replicas $\approx$ 0.5 PB → object store
+ CDN, never the DB.

\section{Interview traps}
\begin{itemize}
  \trap{CAP as ``pick two''; ``NoSQL scales, SQL doesn't'' — most apps fit one
        well-indexed Postgres plus replicas.}
  \trap{A read replica right after a write → the user's change vanishes.}
  \trap{Updating (not deleting) the cache on write; write-back without saying it can lose data.}
  \trap{Averages instead of \textbf{p99}; a single LB/DB = single point of failure.}
\end{itemize}

\section{Remember}
\textbf{Stateless → cache → replicate → shard → queue}, and for each: what
goes stale, what breaks, what number made you do it.

\section{Likely questions}
\begin{enumerate}
  \item L4 vs L7? — connections vs requests; L7 routes by path, terminates TLS.
  \item Cache stampede? — coalesce misses, lock + stale, jittered TTLs.
  \item Post disappears after refresh? — replica lag: read-your-writes from the leader.
\end{enumerate}

\end{multicols}

\noindent{\footnotesize\color{sheetGrey}\textit{Related:} ios-system-design-deep-dives ·
resilience-patterns · event-driven-cqrs-sourcing · hashing-deep (ring) · api-design}

\end{document}
