3.5 Technology deep dives
5.1 The relational model / SQL databases (PostgreSQL, MySQL)
Problem it solves. Store data with regular structure so that arbitrary queries can be answered efficiently without knowing the queries in advance, with the database — not the application — responsible for choosing indexes, join order, and parallelism.
Why wasn't the alternative enough? The hierarchical and network models (1970s–80s) required the application to know the physical access paths; changing a query meant changing the storage layout. Codd's insight was to separate the logical model from the physical access path — that separation is why declarative queries and query optimizers exist at all, and it's why relational survived four waves of challengers.
How it works internally. Tables → rows stored in heap files or a clustered index (Ch 4); B-tree indexes for lookups; a query planner that estimates cardinalities from statistics and picks among nested-loop / hash / merge joins; MVCC for isolation (Ch 8); a write-ahead log for durability and replication (Ch 6).
Deployment. Primary + replicas; connection pooling in front (PgBouncer — Postgres's per-connection process model makes this near-mandatory above a few hundred clients); migrations as versioned, reviewed, forward-only scripts.
Monitoring. Slow query log and pg_stat_statements (total time by normalized query — the single highest-value view); replication lag; connection count vs max_connections; cache hit ratio; transaction ID wraparound / autovacuum progress (Postgres-specific and genuinely dangerous); table and index bloat; lock waits and deadlock rate.
Scaling. Read replicas for read scaling; connection pooling; partitioning by range/hash for very large tables; then sharding at the application layer (Ch 7). Vertical scaling gets you much further than folklore suggests.
Backup. Base backup + WAL archiving for point-in-time recovery (PITR). pg_dump is not a backup strategy for a large production database — it can't do PITR and restore time is measured in hours. Test restores on a schedule; an untested backup is a hypothesis.
What actually breaks in production.
- Long-running transactions block vacuum, bloat accumulates, and eventually the table's performance collapses — or the database refuses writes to prevent XID wraparound.
ALTER TABLEtaking an ACCESS EXCLUSIVE lock behind a queue of waiting queries: the migration itself is fast, but every query behind the lock piles up and the site goes down for the duration. Uselock_timeoutandCONCURRENTLYvariants.- Connection storms — an app tier scaling out past
max_connectionsduring a traffic spike. - The planner flipping to a bad plan after statistics drift, turning a 5 ms query into a 5-minute sequential scan.
- Unbounded
IN (...)lists generated by ORMs. - Replica promoted with data loss because replication was asynchronous (Ch 6).
5.2 Document databases (MongoDB, Couchbase) and JSON columns in relational DBs
Problem it solves. Store self-contained tree-shaped records with good locality, without a rigid up-front schema, in a shape close to the application's objects.
Why wasn't relational enough? Shredding a document-like structure across tables gives cumbersome schemas and unnecessarily complicated application code, plus a multiway join or N queries to reassemble one logical object.
Why isn't it always the answer? Many-to-one and many-to-many relationships don't fit a single document; you cannot address nested items directly; and joins are weak or absent.
How it works internally. A document is stored as one contiguous encoded string (BSON in MongoDB). B-tree indexes can be built on paths inside the document, including fields inside arrays (multikey indexes). Updates typically rewrite the whole document; if the new version doesn't fit in place, it's relocated, which invalidates and updates every index entry pointing at it.
Deployment. Replica sets (primary + secondaries with automatic election); sharded clusters add config servers and routers (mongos).
Monitoring. Document size distribution (the p99 matters far more than the mean); working set vs RAM — once indexes stop fitting in memory, performance falls off a cliff; index hit rate and COLLSCAN counts; replication oplog window (how far a secondary can fall behind before needing a full resync); lock/ticket saturation.
Scaling. Shard on a key with high cardinality and even distribution; avoid monotonically increasing shard keys (timestamps, ObjectIds), which send all writes to one shard.
Backup. Snapshots plus oplog replay for PITR. A mongodump of a sharded cluster is not consistent across shards unless you coordinate it.
What actually breaks.
- Unbounded array growth — the classic "embed comments in the post document" design working fine until a post gets 40,000 comments, at which point every read pulls megabytes and every write rewrites them. This is exactly the "one-to-few, not one-to-many" warning.
- Implicit schema drift. Five years of writes, six generations of shape, and read code full of
if (!user.first_name)branches that nobody dares delete — because nobody knows whether any document still lacks the field. - Application-side joins with no transaction — the two documents fetched are from different points in time.
- Dual-write inconsistency when the relationship is stored on both sides for bidirectional querying.
- Monotonic shard key → one hot shard doing all the writes while the rest idle.
5.3 Graph databases (Neo4j, Memgraph, KùzuDB, Amazon Neptune)
Problem it solves. Queries that traverse a variable, not-known-in-advance number of hops across highly connected heterogeneous data.
Why wasn't relational enough? Every traversed edge is a join with the edges table, and you don't know how many joins you need until you run the query. SQL can express this with WITH RECURSIVE — at roughly 8× the code, with cycle handling and traversal-order control left to you.
Why wasn't a document database enough? Documents model trees. Graphs model arbitrary many-to-many connectivity, in both directions.
How it works internally. Two logical tables (vertices, edges) with indexes on both tail_vertex and head_vertex so traversal works forward and backward. Native graph engines go further with index-free adjacency: a vertex record stores direct pointers to its edge records, so a hop is a pointer dereference rather than an index lookup — this is what makes deep traversals O(hops) rather than O(hops × log n). Query execution then becomes a pattern-matching problem, and the optimizer chooses which end of the pattern to start from (the "forward vs backward" choice in §2.2).
Deployment. Usually a single primary with read replicas — graph databases are notoriously hard to shard, because a good partition of a highly connected graph doesn't exist (min-cut on a social graph is bad by definition). Neptune and Neo4j clusters replicate rather than partition for this reason.
Monitoring. Query traversal depth and expanded-node counts (the real cost driver, not row counts); page-cache hit rate; heap pressure during large traversals; supernode degree distribution.
Scaling. Vertically first, then read replicas. If you truly need to partition, expect cross-partition traversals to dominate cost.
Backup. Full + incremental snapshots; the graph is usually a system of record, so treat it accordingly.
What actually breaks.
- Supernodes. One vertex with millions of edges (a celebrity, a "USA" location node, a shared category) makes every traversal through it explode. Mitigations: edge-type partitioning, degree-aware query planning, or modeling the supernode away.
- Unbounded variable-length patterns —
*0..with no upper bound over a cyclic graph, producing a query that never finishes. Always bound the depth in production. - Cycles causing infinite traversal when the engine doesn't deduplicate visited vertices.
- Sharding attempts that turn every query into a distributed join.
- Schemaless flexibility becoming schemaless chaos — since any vertex can connect to any vertex, nothing stops a bad writer from creating relationships the query code never anticipated.
5.4 GraphQL servers
Problem it solves. Let clients specify exactly the fields their UI needs, so UI changes don't require server changes, and so mobile clients don't over-fetch or make N round trips.
Why wasn't REST enough? Fixed response shapes cause over-fetching (endpoints return more than a screen needs) and under-fetching (a screen needs 4 endpoints, so 4 round trips, or a bespoke endpoint per screen that the backend team must ship).
How it works internally. A schema defines types and the allowed traversals. Execution walks the query tree, calling a resolver per field. The naive implementation is catastrophically N+1: resolving sender for 50 messages calls the user resolver 50 times. The fix is DataLoader-style batching and per-request caching — collect the IDs requested within a tick, issue one batched query, distribute results.
Deployment. A gateway in front of REST/gRPC internal services (which is where most of the operational cost lands — organizations adopting GraphQL often need tooling to convert queries into requests to internal services). Persisted/allow-listed queries in production.
Monitoring. Per-field resolver latency and error rate (not per-endpoint — there is only one endpoint); query depth and complexity score distribution; resolver call counts per request (the N+1 detector); cache hit rate in the batch loaders.
Scaling. Persisted queries (clients send a hash, server has the allow-listed document) both bounds the query space and shrinks requests. Complexity/depth limits. Automatic Persisted Queries + CDN caching for public read traffic.
What actually breaks.
- N+1 resolvers — the default failure mode, and the reason DataLoader exists.
- Denial of service by query complexity — deeply nested or wide queries. The book flags this as the reason the language is deliberately limited; in practice you still need depth limits, complexity scoring, and timeouts, because the schema alone doesn't bound cost.
- Authorization at the wrong layer. Because any field can be reached by many paths, endpoint-level authz doesn't work — authorization must be enforced per field/resolver, and this is where GraphQL security bugs live.
- Rate limiting is hard — one POST to
/graphqlmay be 10 or 10,000 units of work, so request-count limits are meaningless; you need cost-based limits. - Caching is hard — one URL, POST bodies, no HTTP cache semantics for free.
- Schema deprecation — you cannot see who uses a field without field-level usage telemetry.
5.5 Event sourcing platforms (EventStoreDB, MartenDB, Kafka + stream processors)
Problem it solves. Complex business domains where no single representation serves all reads, where intent matters, where auditability is required, and where you want reversibility rather than destructive updates.
Why wasn't a normal database enough? A committed UPDATE/DELETE destroys the prior state and the reason for the change. You cannot re-derive a new view of history you no longer have, and you cannot easily reverse a committed transaction.
Why wasn't plain CDC enough? CDC gives you row-level diffs (Ch 12) — the what, not the why. active=false doesn't tell you whether it was a cancellation, a fraud block, or a data fix.
How it works internally.
- Events appended to a per-aggregate stream with a monotonically increasing sequence number.
- Optimistic concurrency on append: "append these events expecting the stream to be at version N" — this is how command validation stays correct under concurrency, and it's the mechanism most homegrown implementations forget.
- Rebuilding state = fold over the stream. Because folds get slow for long streams, real systems add snapshots (checkpoint state at version N, replay only from there).
- Projections consume the log in order and write read models. Each projection tracks its own checkpoint/offset, which is what makes "delete the view and rebuild" possible.
Deployment. Event store (or Kafka topic with cleanup.policy=compact off — event sourcing needs retention forever, not compaction, unless you're keeping only latest-per-key). Projection workers as separate deployables so they can be rebuilt independently. A schema registry for event types (Ch 5), because old events are never rewritten and must stay readable forever.
Monitoring.
- Projection lag per view (events behind head) — the number that determines whether users see stale data
- Rebuild duration per projection — this is your recovery-time budget, and it grows with history forever
- Append conflict rate (optimistic concurrency retries) — rising rate means an aggregate is a contention hotspot
- Stream length distribution — a stream growing without bound signals a modeling error
- Dead-letter count for events a projection failed on — and remember: a projection is not allowed to reject an event, so any dead letter is a bug, not a business case
Scaling. Partition by aggregate ID (order matters within an aggregate, rarely globally). Snapshot long streams. Run projections independently and in parallel — they don't need to keep pace with each other.
Backup. The log is the backup: views are disposable and rebuildable. So backup discipline concentrates entirely on the event store — and the thing to test is rebuild time, because that's your real RTO.
What actually breaks.
- Non-deterministic projections — the currency-conversion trap, but also
now(), random IDs, and calls to external services inside a projection. Rebuild produces different numbers than the original run and nobody can explain why. - Rebuild time growing past the maintenance window. Fine at 10M events, a weekend outage at 10B.
- Side effects on replay — resending confirmation emails during a rebuild. Requires a strict separation between projections (pure) and reactors/process managers (effectful, with their own idempotency).
- GDPR erasure vs immutability — crypto-shredding works, but destroying the key makes those events unreadable forever, so any future rebuild silently loses them.
- Event schema evolution. You will read 2019 events in 2027. Versioning and upcasting are mandatory, not optional (Ch 5).
- Modeling aggregates too large — one "Conference" stream with a million events means every command replays a million events, or depends entirely on snapshots being healthy.
- Ordering assumptions across streams. Order is guaranteed within a stream. Cross-aggregate ordering in a distributed system is exactly the hard problem of Ch 10.
5.6 DataFrame / array systems (Pandas, Spark, TileDB, NumPy)
Problem it solves. Bridge relational data and the numeric matrices ML algorithms require, with an interactive, incremental "wrangling" workflow.
Why wasn't SQL enough? Feature engineering needs custom code; pivoting into a thousands-of-columns sparse matrix doesn't fit relational storage; and the workflow is exploratory and imperative rather than declarative.
How it works internally. Columnar in-memory representation (Pandas → NumPy arrays; increasingly Apache Arrow as the shared memory format, which is what lets Pandas/Polars/DuckDB/Spark exchange data with zero copies). Operations are vectorized over whole columns. Spark DataFrames are lazy — the chain of transformations builds a logical plan that Catalyst optimizes before any execution, which is why Spark can do predicate pushdown that Pandas cannot.
Deployment. Local single-machine (Pandas, Polars, DuckDB) is right far more often than people assume — see Ch 1's more nodes are not always faster. Distributed (Spark, Dask) when the data genuinely exceeds one machine.
Monitoring. Peak memory vs available (the dominant failure mode); shuffle volume and skew in Spark; task duration skew across partitions; spill-to-disk volume.
Scaling. Prefer a bigger machine before a cluster. Then partition, and above all avoid wide shuffles — a groupBy on a skewed key is the standard Spark cliff.
What actually breaks.
- OOM on a
pivotthat materializes a dense matrix from sparse data — the exact transformation §4 describes, done without a sparse representation. - Silent dtype coercion — a column becomes
objectbecause one row had a string, and everything downstream slows by 100× or computes wrong. - Chained-assignment /
SettingWithCopyWarningin Pandas, where an update lands on a copy and is silently discarded. - Data skew in Spark — 199 tasks finish in 10 s and one runs for 4 hours, because one key holds 60% of the rows.
- Training/serving skew — the one-hot encoding built at training time has different category ordering than at serving time, so the model gets garbage features. The encoder must be persisted with the model, not recomputed.
- Notebook irreproducibility — cells executed out of order, so the DataFrame in memory doesn't correspond to any sequence of code anyone can rerun.