Chapter 01 · Why NoSQL Exists: Workloads, Scale, Flexibility, and Polyglot Persistence
Horizontal Scale, Global Distribution, Flexible Schemas, High Write Rates, and Low-Latency Access
Turn horizontal scale, global distribution, flexible schema, high write rates, and low latency into concrete partitioning, placement, skew, schema-governance, and percentile-latency mechanisms.
Learning outcomes
AtlasMart now has a growth forecast: catalog reads spike during promotions, checkout writes continue during regional traffic shifts, some product records have category-specific attributes, and customers expect low interactive latency. “Use a distributed NoSQL database” is still not a design. Each requirement pulls on a different mechanism: capacity, throughput, partitioning, replication, placement, data shape, and tail latency.
Distinguish scaling up from scaling out, and capacity from throughput.
Separate having replicas in many regions from accepting writes safely in many regions.
Explain why flexible schema moves validation/versioning responsibility rather than removing schema.
Demonstrate how monotonic keys and celebrity tenants create hot partitions even in a many-shard system.
Interpret p50/p95/p99 modeled latency as a distribution and connect it to data placement rather than a database label.
We already decided that database family is not the starting point. This lesson turns five vague requirements—scale, global, flexible, high-write, low-latency—into measurable mechanics that can be placed in the AtlasMart decision notebook later.
1. Scale-up and scale-out solve different constraints
Scale-up gives one node more CPU, memory, faster storage, or network bandwidth. It is operationally simple and can be extremely effective, but eventually reaches hardware, cost, or failure-domain limits. Scale-out distributes data or work across multiple nodes. That increases aggregate capacity and can increase aggregate throughput, but only if the workload is divisible and routing does not concentrate work on one owner.
Capacity is how much state or working set the system can hold. Throughput is how many operations it can complete per unit time under stated latency and durability conditions. Eight nodes do not imply eight times the throughput: replicas consume resources, coordination adds messages, skew creates hotspots, and cross-partition operations can make more nodes participate in a single request.
| Requirement phrase | Mechanism questions that must replace it |
|---|---|
| “We need horizontal scale” | What is the partition key? How many partitions? How are they reassigned? What happens to cross-partition requests? |
| “We need global distribution” | Where are replicas? Which regions may accept writes? What is the inter-region failure behavior? |
| “We need flexible schema” | Who validates required fields and versions? How do old/new readers coexist? How are indexes migrated? |
| “We need high write rate” | Can writes distribute across keys? What is the replication/acknowledgement cost? What storage/compaction work follows? |
| “We need low latency” | Which percentile, from which client region, for what request shape, with what consistency/durability? |
2. Partitioning is the hidden requirement behind most scale-out claims
To scale writes or storage horizontally, the system needs a rule that maps each record or aggregate to a partition (also called a shard in many systems). A routing layer then sends the request to the node or replica set that owns that partition. Range partitioning can preserve locality for ordered scans. Hash partitioning can spread independent keys more uniformly. Neither is universally better.
The intentionally broken case in the lab uses monotonically increasing order numbers and range routing. The current 1,000 writes all land in the newest range, so one shard receives all active write traffic while older shards are idle. Hashing the order ID distributes those writes, but hashing would make an ordered range scan harder. The correct design depends on the dominant access pattern.
A hot shard may sometimes be repaired by changing partition boundaries. A single hot key is harder: if all requests must coordinate on one logical owner, adding unrelated nodes cannot split that key without changing the data model or operation semantics.
3. Global presence is not the same as global write capability
A database may replicate data into several regions so reads can be served near users, yet still direct all writes for a key to one leader region. Another design may accept writes in multiple regions and reconcile or coordinate them. These choices change latency and failure behavior. A local read can be fast but stale; a strongly coordinated write may cross an ocean; an asynchronous write acknowledgement may be fast but expose a larger data-loss or freshness window.
The lab uses declared round-trip-time constants—4 ms local, 82 ms to one remote region, 155 ms to another—to show a distribution. They are model inputs, not measurements. The point is that even when 80% of requests are local, p95 and p99 can be governed by remote traffic. An average would hide that tail.
4. Flexible schema still has a schema
Document stores are often described as “schema-less.” A better phrase is schema-flexible: records can differ in shape, and the database may not require one rigid table definition for every field. The application still needs to know which fields are required, what types and units they use, which versions exist, how indexes treat missing values, and how old readers handle new fields.
AtlasMart catalog products demonstrate a legitimate use. A lamp
has wattage and bulb type; a desk has dimensions and material; a
book has author and ISBN. Embedding category-specific attributes
can avoid a sparse one-size-fits-all table or many subtype
joins. But price_cents, sku, and
schema version are still governed fields. The simulator
validates them explicitly.
5. Run the scale, skew, and placement simulator
Save as atlasmart_scale.py and run with a standard
Python installation. The script performs no networking and
reports no measured database performance.
from collections import Counterimport hashlib, statisticsprint('AtlasMart Lesson 2 — scale-out, placement, skew, and modeled latency')def h(key): return int(hashlib.sha256(key.encode()).hexdigest()[:16], 16)def hash_shard(key, shards=8): return h(key) % shardsdef range_shard(order_number, shards=8, width=1000): return min(shards-1, order_number // width)# Examine the *current* 1,000-write window after seven earlier ranges are full.# Range partitioning preserves locality, but monotonically increasing keys send all new writes to the newest range.current_ids = range(7000, 8000)range_load = Counter(range_shard(i) for i in current_ids)hash_load = Counter(hash_shard(f'o-{i:05d}') for i in current_ids)print('\nCurrent 1,000 writes with monotonic range routing:', sorted(range_load.items()))print('Same writes with hashed order IDs:', sorted(hash_load.items()))print('hashed max/min load ratio:', round(max(hash_load.values())/min(hash_load.values()), 3))# A celebrity tenant remains hot if every request uses the same partition key.tenants = ['tenant-normal-' + str(i) for i in range(7000)] + ['tenant-celebrity'] * 1000tenant_load = Counter(hash_shard(t) for t in tenants)print('\nTenant-key routing with one celebrity tenant:', sorted(tenant_load.items()))print('hottest shard share:', round(max(tenant_load.values())/len(tenants), 3))# Modeled network latency. These are declared model constants, not benchmark results.region_rtt_ms = {('eu','eu'):4, ('eu','us'):82, ('eu','ap'):155}requests = [('eu','eu')]*80 + [('eu','us')]*15 + [('eu','ap')]*5latencies = [region_rtt_ms[x] for x in requests]latencies.sort()def pct(xs, p): return xs[min(len(xs)-1, int(round((len(xs)-1)*p)))]print('\nModeled RTT distribution (not a measured database benchmark):')print('p50=', pct(latencies,.50), 'ms p95=', pct(latencies,.95), 'ms p99=', pct(latencies,.99), 'ms')# Flexible shape requires explicit schema/version handling.docs = [ {'schema_version':1,'sku':'p-1','name':'Lamp','price_cents':2500}, {'schema_version':2,'sku':'p-2','name':'Desk','price_cents':18900,'dimensions_cm':[120,60,75]},]def validate(d): required={'schema_version', 'sku','name','price_cents'} missing=required-set(d) if missing: raise ValueError(f'missing fields: {sorted(missing)}') if d['price_cents'] < 0: raise ValueError('negative price')for d in docs: validate(d)print('\nFlexible documents validated:', len(docs))print('Lesson: scale-out changes routing/failure problems; flexible shape moves schema responsibility rather than eliminating it.')
Verified deterministic output
AtlasMart Lesson 2 — scale-out, placement, skew, and modeled latencyCurrent 1,000 writes with monotonic range routing: [(7, 1000)]Same writes with hashed order IDs: [(0, 111), (1, 117), (2, 117), (3, 128), (4, 135), (5, 128), (6, 127), (7, 137)]hashed max/min load ratio: 1.234Tenant-key routing with one celebrity tenant: [(0, 872), (1, 859), (2, 882), (3, 868), (4, 895), (5, 879), (6, 898), (7, 1847)]hottest shard share: 0.231Modeled RTT distribution (not a measured database benchmark):p50= 4 ms p95= 82 ms p99= 155 msFlexible documents validated: 2Lesson: scale-out changes routing/failure problems; flexible shape moves schema responsibility rather than eliminating it.
The first block is a controlled failure: all current range-routed writes target shard 7. The hashed representation distributes the same 1,000 keys across eight shards. The celebrity tenant then shows why hashing is not magic: every request for the same tenant key still hashes to the same shard, lifting that shard to roughly 23% of the total requests in this synthetic workload.
6. High write rate depends on what happens after acknowledgement
A front-end write is not finished merely because a client received “OK.” The system may append to a write-ahead log, replicate to peers, update memory structures, later flush immutable files, maintain secondary indexes, compact old files, or repair replicas. A write path can therefore look fast at low load while accumulating background debt. Chapter 09 will make WAL, LSM trees, SSTables, compaction, and amplification observable.
For Chapter 01, record the questions: how many replicas must acknowledge before success? Can two regions write the same key? Can retries create duplicates? What is the durable recovery point if a node or region fails after acknowledgement? Does one secondary index turn one logical write into several physical writes? A “high-write database” claim is meaningless without those settings.
7. Production judgment and acceptance criteria
Scale-out is appropriate when a measured workload exceeds one failure domain or one node’s practical capacity/throughput, or when placement requirements demand several locations. It is not free insurance. Every new partition introduces routing metadata and rebalancing; every replica consumes storage and bandwidth; every writable region creates coordination or conflict semantics; every flexible shape needs compatibility rules.
For AtlasMart, the next architecture review should refuse requirements like “global and fast.” Replace them with: p95/p99 targets by operation and client region; peak read/write rates; record/key-size distributions; hot-key estimates; replica count; acknowledgement/durability policy; allowed stale-read window; region-failure behavior; RPO/RTO; and cost headroom during rebalancing or one-node loss.
Verification checklist
- The simulator clearly distinguishes range and hash placement.
- The range-hotspot example uses the current write window, not a misleading all-time average.
- The celebrity key remains hot after hashing, demonstrating that node count alone is insufficient.
- Latency values are labeled as model constants, not measured product results.
- Both document versions pass explicit validation, reinforcing that flexible shape is governed shape.
Check your understanding
- Why can eight shards still behave like one overloaded shard?
- What tradeoff does hash partitioning make compared with range partitioning?
- Why does placing a replica in a region not prove that writes are local there?
- What must a flexible-schema application version or validate?
- Why are p95/p99 often more useful than an average for user-facing latency?
Review the answers
Skew can concentrate traffic on one partition or one hot key. Aggregate cluster capacity does not help a request path that has one mandatory owner.
Hashing tends to distribute unrelated keys but destroys natural key ordering/locality, so range scans may need fan-out or another index.
A replica may be read-only or may forward writes to a leader elsewhere. Multi-region write acceptance has separate consistency/conflict semantics.
Required fields, types, units, compatibility rules, schema versions, defaults, index behavior, and migration logic still need ownership even without rigid table DDL.
Tails show the slower fraction users experience and expose remote routes, queueing, pauses, or overloaded partitions that an average can hide.
Authoritative references
- Dynamo: Amazon’s Highly Available Key-value Store — Primary source for partitioning, replication, and availability-oriented key-value architecture.
- Bigtable: A Distributed Storage System for Structured Data — Primary source for large-scale distributed structured data and tablet/range-oriented design ideas.
- Google SRE: Service Level Objectives — Official SRE guidance on measurable SLIs/SLOs and percentile-oriented service objectives.
- Python 3.14.7 release — Official current stable Python feature-release reference at chapter-generation time.