1.5 Technology deep dives
For each: what problem, why wasn't the alternative enough, how it works internally, deployment, monitoring, scaling, backup, what breaks in production.
For each: what problem, why wasn't the alternative enough, how it works internally, deployment, monitoring, scaling, backup, what breaks in production. Book content is the concepts; the ops sections are the practitioner layer the book leaves to you.
5.1 Object storage (Amazon S3 / Azure Blob / GCS / Cloudflare R2)
Problem it solves. Store an effectively unbounded amount of large files durably, without any single machine's disk capacity or failure being your problem.
Why wasn't a filesystem / RAID / NAS enough?
- A filesystem is bounded by one machine's disks. RAID survives a disk failure, not a machine or datacenter failure.
- NAS/SAN scale vertically and require capacity planning; you must decide the size in advance and pay for idle space.
- Object storage replaces capacity planning with metered billing and replaces machine-awareness with an API.
How it works internally.
- Flat key → blob namespace per bucket (the "directories" are a UI fiction over
/-containing keys). - Objects are split into chunks, chunks are replicated (or erasure-coded, e.g. Reed–Solomon k-of-n) across independent failure domains — different disks, machines, racks, and availability zones. Erasure coding gives similar durability to 3× replication at ~1.5× storage overhead.
- A separate metadata/index service maps keys → chunk locations. This index is the real scaling challenge, and it's why per-key operations are strongly consistent while listing is often weaker.
- Immutable objects: a PUT to an existing key writes a new version and flips a pointer. There is no in-place partial write — this is the single most important behavioral fact for system design on top of S3.
- Consistency: S3 has offered strong read-after-write consistency for PUTs and DELETEs since Dec 2020; before that it was eventually consistent, and enormous amounts of old design advice assumes the old model.
Deployment. No deployment — it is the substrate. Self-hosted equivalents: MinIO, Ceph RADOS Gateway, SeaweedFS. What you do deploy is policy: bucket per environment/tenant, IAM/bucket policies, encryption (SSE-S3 vs SSE-KMS), versioning on, lifecycle rules (transition to infrequent-access/Glacier at N days, expire noncurrent versions at M days), and Block Public Access at the account level.
Monitoring.
- Request rates and 4xx/5xx split — a rising 503
SlowDownrate means you're hitting per-prefix request limits - First-byte and total latency percentiles (p50/p99), not averages
- Bytes stored by storage class and by prefix (this is your bill)
- Replication lag if using cross-region replication
- Access logs / CloudTrail data events for who read what — required for most compliance stories
Scaling. Effectively unlimited capacity. The real limit is request rate per prefix (S3: ~3,500 PUT/COPY/POST/DELETE and ~5,500 GET/HEAD per second per partitioned prefix). You scale by spreading keys across many prefixes — historically by putting a hash at the start of the key. Read scaling beyond that: CloudFront/CDN in front, or S3 Transfer Acceleration for long-haul uploads.
Backup. Durability (11 nines) is not backup — it protects against disk failure, not against you. You need: versioning (undo overwrite/delete), MFA Delete or an Object Lock retention policy (undo malicious delete), cross-region/cross-account replication (undo account compromise or region loss). The account boundary matters: a backup in the same account that an attacker owns is not a backup.
What actually breaks in production.
- Eventual-consistency assumptions in old code — list-after-write still isn't strongly consistent in every provider/operation; jobs that "list the directory then process" silently miss new files.
- Hot prefix throttling — a batch job writing
2026-08-19/part-0000…writes every key to one prefix and gets 503-throttled. - Cost explosions from small objects. Object stores are for hundreds-of-KB-to-GB files. Millions of 2 KB objects means you pay per-request costs that dwarf storage, and any framework listing them takes forever. This is exactly why cloud databases pack many small values into large blocks before writing to the object store.
- Lifecycle rules that delete data someone still depended on — because the dependency was undocumented.
- Silent public exposure through a permissive bucket policy or a pre-signed URL with a year-long expiry.
- Multipart upload leakage — aborted multipart uploads keep consuming (billed) storage forever unless you add a lifecycle rule to abort incomplete uploads.
5.2 Data warehouse (Snowflake / BigQuery / Redshift / Synapse)
Problem it solves. One place where analysts can join data from every operational system and run arbitrary, expensive queries without touching production.
Why wasn't running analytics on the OLTP DB enough? The four reasons in §1.4: data silos, wrong schema, performance interference, network/compliance isolation.
How it works internally. (Full treatment in Ch 4; here's the shape.)
- Columnar storage — each column stored contiguously, so a query touching 2 of 200 columns reads 1% of the bytes.
- Compression per column (run-length, dictionary, bit-packing) — homogeneous data compresses hugely.
- Vectorized execution — operate on batches of column values in tight loops, exploiting SIMD, instead of row-at-a-time.
- Disaggregated storage/compute — Snowflake stores in S3 and spins up independent "virtual warehouses" (compute clusters) per workload, which is how one team's heavy query stops affecting another's. BigQuery separates storage (Colossus) from compute (Dremel) with the Jupiter network in between.
- Micro-partitions + zone maps — per-block min/max metadata lets the engine skip blocks entirely without reading them (data skipping / partition pruning).
Deployment. Warehouse itself is managed. What you deploy is the modeling layer: dbt (or equivalent) for transformations under version control, with CI running tests on models; separate dev/staging/prod databases; scheduled orchestration (Airflow/Dagster) for pipeline DAGs.
Monitoring.
- Cost per query and per warehouse — the #1 warehouse ops metric; bytes scanned (BigQuery) or credit-seconds (Snowflake)
- Query queueing / concurrency saturation
- Freshness/SLA per table — "how stale is this table?" is the metric analysts actually feel
- Pipeline failure and retry counts, and row-count/schema drift tests on each model
- Spill to disk (a query exceeding memory), which is the signal to resize the warehouse or fix the query
Scaling. Scale compute independently of data: resize a virtual warehouse (bigger machines) for a single heavy query, or add clusters (multi-cluster auto-scale) for concurrency. Isolate workloads on separate warehouses — ELT, BI, and data science shouldn't share compute.
Backup. Managed platforms give time travel (Snowflake: 1–90 days, then Fail-safe 7 days; BigQuery: 7-day time travel + snapshots). Time travel is not disaster recovery — for that you need cross-region replication and, for regulated data, exports to your own object storage. The real backup for a warehouse is the ability to rebuild it from the lake/source systems, because the warehouse is derived data. Test that rebuild.
What actually breaks in production.
- A single analyst's exploratory query costs thousands of dollars (an unfiltered scan of a petabyte table). Fix: partitioning + clustering + required partition filters + per-user byte quotas.
- Silent pipeline failure — a job fails, no one notices, dashboards show yesterday's numbers as if they were today's. Freshness alerting, not just job-failure alerting.
- Schema drift upstream — a source system renames a column; the ETL silently nulls it, and a KPI quietly drops to zero.
- Timezone and late-arriving-data bugs — the most common source of "the numbers don't match" between two dashboards.
- Two teams computing "revenue" differently — not a technical failure, but the one that destroys trust in the warehouse fastest. This is what the semantic/metrics layer exists to prevent.
5.3 Data lake (S3/ADLS + Parquet + a table format)
Problem it solves. Store any data — including unstructured — cheaply, in raw form, so each consumer can transform it their own way (the sushi principle).
Why wasn't a warehouse enough? Feature engineering, NLP, and computer-vision workflows need custom code over arbitrary formats; data scientists prefer Pandas/scikit-learn/R/Spark to SQL; and relational storage is more expensive than object storage.
How it works internally. Files on an object store, typically Parquet (columnar, compressed, with per-row-group statistics) or Avro (row-oriented, great for streaming/CDC because it evolves schemas cleanly). A table format (Apache Iceberg, Delta Lake, Apache Hudi) layers over the files a metadata log giving: atomic commits, snapshot isolation, schema evolution, partition evolution, and time travel — this is what turned "a lake" into "a lakehouse" and made it queryable safely by multiple writers.
Deployment. Bucket layout by zone (raw/, staging/, curated/), partitioned by ingestion date; a catalog (Glue/Hive Metastore/Unity/Iceberg REST catalog) as the single source of table metadata; query engines (Athena/Trino/Spark/DuckDB) pointed at the catalog.
Monitoring. Small-file count per table (the canonical lake health metric), compaction job success, partition skew, catalog/metadata staleness, bytes scanned per query, and orphan-file accumulation.
Scaling. Storage scales for free. Query scaling = partition pruning + file sizing (aim for ~128 MB–1 GB Parquet files) + Z-ordering/clustering on high-selectivity predicates.
Backup. Object versioning + cross-region replication for the files; the catalog is the fragile part — back up the metastore/catalog database, because losing it turns a well-organized lake back into an undifferentiated pile of files.
What actually breaks.
- The small-file problem — a streaming job writing every minute produces millions of tiny Parquet files; query planning time exceeds query time. Requires scheduled compaction.
- The lake becomes a swamp — no catalog, no ownership, no schema, nobody knows which of the 40
users_final_v3datasets is real. - Concurrent writers corrupting a table without a table format's atomic commits (partially-written partitions read as real data).
- Schema-on-read mismatches — one file has
user_idas string, the next as int; the query engine fails, or worse, silently coerces. - GDPR erasure — deleting one person's rows from immutable Parquet requires rewriting whole files. Table formats' row-level deletes exist precisely for this, and are slow and easy to forget to run against all derived copies.
5.4 ETL/ELT connectors (Fivetran / Airbyte / Singer)
Problem it solves. Getting data out of SaaS products whose database you can never access — only a rate-limited, paginated, idiosyncratic API.
Why not write it yourself? You can — once. The cost is maintenance: every SaaS vendor changes their API, their pagination, and their rate limits independently, and each of your 40 connectors breaks on a different Tuesday.
How it works internally. Per-source connector with an incremental cursor (updated-at watermark, or a change/log endpoint), a normalization step into a target schema, and idempotent upserts into the destination keyed on primary key. Good connectors checkpoint the cursor so a failed sync resumes rather than restarting.
Monitoring. Sync success/failure per connector, rows synced vs expected, sync duration trend, API rate-limit consumption, and replication lag per table (the number analysts care about).
Scaling. Bounded by the source API, not by you. Scale = more frequent incremental syncs, never more parallel full-refreshes.
What actually breaks. Silent partial syncs after an API deprecation; hard deletes at the source that a watermark-based incremental sync can never see (you need soft deletes or periodic full refreshes); a full re-sync triggered by a schema change consuming the month's API quota in an hour; and PII arriving in the warehouse that nobody realized the SaaS was returning.
5.5 Real-time OLAP (ClickHouse / Druid / Pinot)
Problem it solves. Analytical queries (aggregates over millions/billions of rows) at interactive latency inside a user-facing product — the dashboard your customers see, not the one your analysts see.
Why wasn't a data warehouse enough? Warehouses batch-ingest and optimize for throughput; query latency of seconds-to-minutes is fine for an analyst and fatal in a product. Why wasn't Postgres enough? It's row-oriented; scanning a billion rows to compute a sum reads every column of every row.
Why wasn't a precomputed rollup enough? It is, until users want a filter dimension you didn't precompute. These systems buy you arbitrary slicing.
How it works internally (ClickHouse as the concrete case).
- MergeTree family: data is written as immutable parts, sorted by the table's
ORDER BYkey; a background process merges parts into larger ones (an LSM-tree lineage — see Ch 4). - Columnar storage with per-column codecs (
Delta,DoubleDelta,Gorilla,LZ4,ZSTD); sparse primary index (one index entry per ~8,192-row granule) so the index fits in memory even for trillion-row tables. - Vectorized execution over 65,536-row blocks; heavy SIMD use.
- Materialized views that trigger on insert, plus
AggregatingMergeTree/SummingMergeTreeto maintain rollups incrementally. - Distribution:
Distributedtable fans a query out to shards; each shard replicated (ReplicatedMergeTree coordinates via Keeper/ZooKeeper).
Deployment. Shards × replicas; a coordination service (ClickHouse Keeper) with an odd number of nodes; ingestion via batched inserts (never row-at-a-time) or Kafka table engine. Separate the ingest path from the query path.
Monitoring. Parts-per-partition (the health metric — too many means merges are losing), merge queue depth and background pool saturation, insert batch size, max_memory_usage rejections, replication queue depth, mutation progress, and query p99 by query type.
Scaling. Vertical first — these engines exploit big machines very well. Then shard on a key with even distribution, and add replicas for read concurrency. Cluster reads and writes onto different replicas if query load is spiky.
Backup. BACKUP/RESTORE to object storage, or filesystem-level snapshots of parts plus schema DDL in version control. Replication is not backup — a bad ALTER DELETE replicates instantly.
What actually breaks.
- "Too many parts" errors — caused by frequent small inserts. The single most common ClickHouse production incident, and it's an application bug (not batching), not a database bug.
- Choosing a bad
ORDER BY— irreversible without a full rewrite, and it determines every query's performance. - Memory-blown queries — a
GROUP BYon a high-cardinality column OOMs the server, taking down queries for everyone. - Distributed queries with no
GLOBAL IN— an innocuous subquery is re-executed per shard and melts the cluster. - Mutations (
ALTER UPDATE/DELETE) are asynchronous rewrites of whole parts; people treat them like OLTP updates and are surprised when a "delete one row" rewrites 200 GB.
5.6 HTAP (SingleStore / TiDB / Spanner-style hybrids)
Problem it solves. One application needing both low-latency single-record reads/writes and large scans — fraud detection is the archetype, where you must score a transaction now against aggregate history.
Why wasn't OLTP + ETL + OLAP enough? ETL latency. If the decision must be made in the request path, a warehouse refreshed every 15 minutes is useless.
How it works internally. Typically a row store for recent/hot data + a column store for historical data, with a background process converting row-format to column-format, and a query planner that can read both and union the results. This is the "internally two systems behind one interface" the book warns about.
What actually breaks. The seam between the two engines: queries that unexpectedly hit the row store scan and blow up; conversion lag making "analytics" silently exclude the last N minutes; and the operational reality that you now have one system with two very different tuning models and one team that understands neither fully.
5.7 Kubernetes (as the microservices substrate)
Problem it solves. Every microservice independently needs: releases, resource allocation, log collection, health monitoring, and alerting. K8s provides that foundation once instead of N times.
Why wasn't a VM per service with config management enough? Bin-packing (one VM per service wastes capacity), rollout speed, and self-healing. Why wasn't a PaaS enough? Less control over networking, storage, and scheduling.
How it works internally. A declarative control loop: you write desired state to etcd via the API server; controllers continuously reconcile actual state toward desired; the scheduler binds Pods to Nodes by resource requests and constraints; kubelet on each node starts containers; kube-proxy/CNI provide service networking. Everything is a reconciliation loop — that's the whole design.
Monitoring. Pod restart counts and CrashLoopBackOff, OOMKilled counts, CPU throttling (a container hitting its CPU limit is throttled, not killed — and this is the silent latency killer), pending pods (= no node fits the request), node pressure conditions, etcd latency and DB size, and control-plane API latency.
Backup. etcd snapshots are the cluster backup; plus your manifests in Git (GitOps) so the cluster is reproducible. PersistentVolumes need their own backup (Velero, or CSI volume snapshots) — a common and painful gap.
What actually breaks.
- Requests set too low → nodes oversubscribed → everything throttles simultaneously under load.
- Missing liveness/readiness distinction — a liveness probe that fails under load restarts healthy-but-busy pods, converting a slowdown into an outage.
- PodDisruptionBudgets missing — a node drain during an upgrade takes all replicas of a service at once.
- Stateful workloads (databases) on K8s without understanding the storage layer — the classic way to lose data.
- DNS, always. CoreDNS saturation shows up as random, inexplicable timeouts across unrelated services.
5.8 Distributed tracing (OpenTelemetry / Jaeger / Zipkin)
Problem it solves. In a distributed system, "the system is slow" has no local answer. Tracing lets you ask which call, in which service, for which operation, took how long.
Why weren't logs and metrics enough? Metrics tell you that p99 rose; logs tell you what one service did. Neither reconstructs a single request's path across service boundaries, which is where the latency actually hides.
How it works internally. A trace = a tree of spans. Each span has a trace ID, span ID, parent span ID, timestamps, and attributes. Context propagation injects the trace ID into outgoing request headers (W3C traceparent), so downstream services attach their spans to the same trace. Spans are batched, exported to a collector, and sampled — because storing every span at scale is prohibitive. Two sampling models: head-based (decide at the root, cheap, may miss rare slow traces) and tail-based (buffer the whole trace, then keep the slow/erroring ones — much more useful, much more expensive, and requires all spans of a trace to reach the same collector).
Monitoring the monitoring. Span export failure/drop rate, collector queue saturation, and cardinality of span attributes — an attribute containing a user ID or a raw URL with IDs in it will destroy your tracing backend's cost model.
Scaling. Collector as a horizontally scaled deployment (with a load-balancing exporter in front when doing tail sampling), aggressive sampling of high-volume/low-value routes, and strict attribute allow-lists.
What actually breaks. Broken context propagation across an async boundary (a queue, a thread pool) — traces silently truncate at exactly the place you needed to see. Instrumenting only the HTTP layer, so all latency appears as "the database call" with no detail. And PII in span attributes, which turns your observability stack into a compliance liability.