EE 547 - Unit 5
Dr. Brandon Franzke
Fall 2026
Document, Wide-Column, and Graph Stores
Primary: the one server that accepts writes
Replica: a second server, with its own disk and memory
Routing: two addresses, one for writes and one for reads
Read load: the departures board, reports, any table read far more than it is written
Loss of the primary: a current copy already on another server
primary> UPDATE flights SET seats_remaining = 0
WHERE flight_id = 7;
UPDATE 1
replica> SELECT seats_remaining FROM flights
WHERE flight_id = 7;
1
replica> [same query, a moment later]
0
Replication lag: the time from the commit on the primary until the replica has applied the same change
Stale rows: both answers correct for the replica’s data at that moment
Write-ahead log
Shipping
Lag: the records the replica has not yet applied
Commit
Two disks: every committed write on two servers
Round trip: added to every write
Unreachable replica: every commit waits until the replica is reachable again or is removed from the pair
On AWS: RDS Multi-AZ, a standby in a second availability zone
Commit
No added wait: a write completes as on one server
Lag: each replica serves reads as of its last applied record
Loss of the primary: records not yet sent are lost
On AWS: a read replica
Email change
One reader notices: the client that wrote the new value, the only reader that has seen it
Routing fixes
Which reads follow a client’s own write is visible only to the application
Promotion
Data on the new primary
Lost records: with an asynchronous replica, the records not yet sent at the failure
Read capacity: each replica serves reads from its own disk and memory
Continuity: a lost primary is replaced by a copy, with no restore from backup
Distance: a replica in Frankfurt serves Frankfurt’s reads
Write capacity: one server accepts every write
Storage: every copy holds all the data
Shard: a server holding part of the rows of every table
Shard key: the column a row’s shard is computed from
flight_id for flights and for their bookingsRouting: each read or write sent to the shard holding its key, by the application or by a router in front of the shards
Capacity: the shards together hold more rows than one disk and accept more writes a second than one server
By the shard key: WHERE flight_id = 7
Without the shard key: WHERE departure_time > now()
flights rowsWith flight_id as the shard key, the departures board, the most frequent read, goes to every shard
Hash function: maps a key to a number in a fixed range
hash("N12345") = 2 303 133 630 mod 4 = 2
hash("N12346") = 273 569 284 mod 4 = 0
hash("N12347") = 1 732 863 634 mod 4 = 2
(same shard as N12345, by chance)
Placement: hash(key) mod N is the shard number, for N shards
Even spread: on average each shard holds the same number of rows and accepts the same share of the writes, for any set of keys
One shard per key: a read or write for a key goes to exactly one server
Neighbors apart: N12345 and N12346 on different shards
Placement: each shard holds one contiguous slice of the keys in sorted order
Range queries: go to the one or two shards whose slices cover the range
Uneven spread: the number of rows in a slice depends on the keys
Hot shard: with a timestamp as the key, every row written this minute goes to the shard holding the current slice
Hash when reads are by one key and writes should spread; range when reads cover a range of keys and the key is not the time of arrival
With hash(key) mod N, a change in N changes the shard of most rows
Partition: a fixed slice of the hash range
Partition count: set far above the number of shards, hundreds to tens of thousands (Redis Cluster: 16,384), or partitions split as they grow (MongoDB chunks, DynamoDB partitions)
Shard added: the new shard is given whole partitions, one from each existing shard in turn
Shard removed: its partitions are given to the remaining shards
Ring: the same scheme drawn on a circle (Cassandra’s form)
The hash of a key never changes, only the partition table
Join on one shard: a flight and its bookings placed by the same flight_id
Join across shards: bookings placed by passenger_id
Two-phase commit: the seat count on one shard and the new booking row on another
Failure between the phases: a shard that has returned yes holds its locks until the coordinator’s decision arrives
A commit is one log write on one shard and two round trips across two
Reads by one key: placement by that key, flight_id for flights and their bookings
Reads over a range: range placement on that column, never on a column ordered by time of arrival
Compound key: hash on one column, rows sorted by a second column inside the shard
Departures board: placed by origin, then departure_time
| Read | Shard key | Shards |
|---|---|---|
| Seat count, flight 7 | flight_id |
one |
| Bookings for flight 7 | flight_id of the booking |
one, the flight’s |
| One passenger’s bookings | flight_id of the booking |
every shard |
| One aircraft’s reports, last hour | aircraft_id, then time |
one |
| Departures from LAX, next two hours | flight_id |
every shard |
| Departures from LAX, next two hours | origin, then departure_time |
one |
A read that filters on a column other than the shard key goes to every shard or to a second copy of the table placed by that column
Network partition: the links between two groups of servers fail
app (zone a)> UPDATE flights ... -- primary, zone b
connection to server at "db-b.internal"
(10.0.2.14), port 5432 failed:
Connection timed out
app (zone a)> SELECT ... -- replica, zone a
seats_remaining
1
Causes
Duration: seconds for a reroute, minutes for a switch restart, hours for a repair
Each side
Both sides hold copies of the same rows, and both are receiving writes
Consistency: every read returns the most recent write, from whichever server answers it
Availability: every request to a working server gets an answer that is not an error
Partition tolerance: the database continues operating while its servers are split into groups that cannot communicate
Pick two: a database has at most two of the three
Partition tolerance: required, since links fail
A write arrives at a server cut off from the other copies
Without a partition
Per operation: one database answers for some operations and refuses for others
Majority rule: a write completes only when more than half of the copies confirm it, two of three
During the partition
After the link returns: the minority copy receives the writes it missed, with nothing to reconcile
Copies and losses
During the partition
seats = 0After the link returns: one of three rules resolves the disagreement
Eventual consistency: the copies agree once every copy has applied every write and each conflict is resolved
ACID: the relational transaction’s guarantees, per transaction
Single primary: one server applies every write
BASE: the design that answers on both sides of a partition
Per-operation level
ACID refuses during a partition and BASE reconciles after it
Clock skew: synchronized clocks on two servers differ by milliseconds
Version vector: a count of the writes applied from each server, stored with the row, such as {A: 3, B: 5}
Read result
Three counts
Overlap: with \(W + R > N\), every read set shares at least one copy with the last write set, and that copy has the write
Settings
Cassandra sets ONE, QUORUM or ALL per operation, and DynamoDB a strongly consistent read
Consistency level: W and R set per operation by the cost of a stale or a refused answer, on the same three copies
Seat sale: majority for the write and the read
Departures board: one copy for the write and the read
Fixed level: a database with one level applies the same choice to every operation
Under a partition, the seat sale refuses on the small side and the board continues
Key-value store: key and value pairs in memory on one or more servers, beside the database
Key: a string the application chooses, and the only thing indexed
Value: bytes, opaque to the store
Instances
Stored values: the login session, the cached departures board, a request count, the standby list
import redis, json r = redis.Redis(host="cache.internal") r.set("session:7c1e", json.dumps({"passenger_id": 42})) r.get("session:7c1e") # b'{"passenger_id": 42}' r.get("session:0000") # None
Get: one key in, the bytes out, or nothing
Round trip: one request to the server the key hashes to
Serialization: object to bytes and back, in the application
r.set("flight:7:seats", 180) r.get("flight:7:seats") # b'180' r.delete("flight:7:seats") # gone after 30 min; expire() resets it r.set("session:7c1e", data, ex=1800) r.expire("session:7c1e", 1800) # 38, in one step r.incr("rate:42:14:05") # True once, then None; released at 30 s r.set("lock:board:LAX", "w3", nx=True, ex=30) # one round trip for both r.mget(["flight:7:seats", "flight:12:seats"])
Set, get, delete: a whole value at a time, each write replacing the bytes
Expiry: a lifetime in seconds on the key
Atomic operations: increment; set-if-absent
Many keys: one round trip per key
There is no query language beyond these operations
Key pattern: entity, identifier, attribute, joined by a separator
session:<token> login state, one passenger
flight:<flight_id>:seats seat count, a counter
board:<airport> departures board, one value
standby:<flight_id> standby list, in order
rate:<passenger_id>:<minute> requests this minute
Lookups with a key: the session for token 7c1e; the board for LAX
Lookups without one: all sessions of passenger 42
Second key: the set stored under a key of its own, written on every change
r.sadd("passenger:42:sessions", "7c1e") # at login r.smembers("passenger:42:sessions") # {b'7c1e', b'a91b'}
Planned lookups: a key-value store answers only the lookups its keys were built for
SCAN walks every key on every server and returns each valueSCAN is a maintenance operation, outside the request pathRelational schema, by key
Write path
Persistence settings* (*Redis defaults; other stores name the same choices differently)
Replication: the replica receives each write after the acknowledgement, as an asynchronous database replica does
Loss at a crash
Sessions, cached values and counters can be lost and rebuilt without harm
def departures(airport): key = f"board:{airport}" hit = r.get(key) if hit is not None: return json.loads(hit) rows = db.query( "SELECT ... FROM flights " "WHERE origin = %s " "AND departure_time > now() ...", airport) r.set(key, json.dumps(rows), ex=30) return rows
Hit: the value is in the cache; one round trip, no query
Miss: the key is absent or expired
Hit rate: the fraction of reads answered from the cache
Two copies: the row in the database and its copy in the cache
Lifetime: the copy is served until it expires, regardless of writes to the database
Delete on write: the cached key deleted with the row’s write, and the next read refills it
db.execute( "UPDATE flights SET departure_time = %s " "WHERE flight_id = 7", t) r.delete("board:LAX")
Eviction: at full memory, the cache removes keys, least recently used first
Hot key: one key read far more than the rest
The board for a hub airport, on every display and every phone in the terminal
Hashed to one server
Adding servers spreads the other keys, not this one
That server’s link and CPU are the ceiling for this key
At expiry: the key is gone until the first miss has run the query and written the result
Claiming the refill: set-if-absent on a lock key succeeds for one request
The same arithmetic applies to every popular key at every expiry
Rate limit: 100 requests a minute per passenger
# read, check, write: # two requests can both read 99 n = int(r.get(f"rate:{pid}:{minute}") or 0) if n < 100: r.set(f"rate:{pid}:{minute}", n + 1) allow() # both write 100: 101 admitted # one operation: the server increments # and returns the new value n = r.incr(f"rate:{pid}:{minute}") # 100, then 101 if n == 1: r.expire(f"rate:{pid}:{minute}", 60) if n <= 100: allow()
Interleaving: another client’s write between this client’s read and write
Atomic step: increment, or set-if-absent, performed by the server as one indivisible operation
Guarded update: UPDATE flights SET seats_remaining = seats_remaining - 1 WHERE seats_remaining > 0, read, check and write in one statement
INCR: the same shape, by keySeveral keys at once
MULTI … EXEC: a group of commands run without interleavingWATCH: the group is abandoned if a watched key changed since it was readLock: SET lock:flight:7 <worker> NX EX 30, the claim and its expiry in one operation
True; the rest get None and wait or move onAtomicity: one operation, or one group on one server, never across servers
Standby list for flight 7: passengers ordered by priority; read as “the next three” and “where am I”
# score: priority r.zadd("standby:7", {"p42": 2.0, "p17": 1.0, "p88": 2.5}) r.zrange("standby:7", 0, 2) # [b'p17', b'p42', b'p88'] r.zrank("standby:7", "p88") # 2 r.zrangebyscore("standby:7", 1.0, 2.0) r.zrem("standby:7", "p17") # seat assigned
Score: a number per member, with members stored in score order
Operations on the server: insert, rank, range by rank, range by score
Other value types: hash (field and value pairs in one key), list, set
Key-value fit: data the application names by a key it already has, reads or writes whole, and can lose or rebuild
Departures board: fits, as a copy
Seat count: does not fit
Facts are stored in the tables, copies in the key-value store
bookings = db.query( "SELECT * FROM bookings " "WHERE passenger_id = %s", pid) for b in bookings: # N rows b.flight = db.query( # one query each "SELECT * FROM flights " "WHERE flight_id = %s", b.flight_id)
N + 1: one query for the list, then one per row
Join: the same rows in one query
SELECT b.*, f.origin, f.destination, f.departure_time FROM bookings b JOIN flights f ON f.flight_id = b.flight_id WHERE b.passenger_id = %s;
N + 1 appears wherever the application follows identifiers itself
Passenger’s trips: one read of the passenger’s upcoming flights with seat, flight and aircraft
passengersbookingsflightsaircraftNormalized tables: each fact once, in its own table
Access pattern: a read or write the application makes, with the facts it needs together and how often it runs
Normalization: each fact in one place
Denormalization: the fact copied into every document read with it
Read-heavy facts: reads far outnumber writes
Entity-first design in a store without joins: one collection per entity, the join done by the application
-- normalized: five tables, six joins, per read SELECT b.booking_id, p.name, p.email, f.flight_number, f.departure_time, o.city AS from_city, d.city AS to_city, ac.model FROM bookings b JOIN passengers p ON p.passenger_id = b.passenger_id JOIN flights f ON f.flight_id = b.flight_id JOIN airports o ON o.airport_code = f.origin JOIN airports d ON d.airport_code = f.destination JOIN aircraft ac ON ac.aircraft_id = f.aircraft_id WHERE b.booking_id = 9001;
{ "booking_id": 9001, "passenger": { "name": "A. Rivera", "email": "a.rivera@..." }, "flight": { "number": "UA 455", "departure_time": "2026-10-09T14:05", "from_city": "Los Angeles", "to_city": "San Francisco", "aircraft": "Airbus A320" } }
{ "passenger_id": 42, "name": "A. Rivera", "bookings": [ { "booking_id": 9001, "seat": "14C", "fare_class": "Y", "flight": { "flight_id": 7, "number": "UA 455", "origin": "LAX", "destination": "SFO", "departure_time": "2026-10-09T14:05", "aircraft": "Airbus A320" } }, { "booking_id": 9144, "seat": "2A", "fare_class": "J", "flight": { "flight_id": 12, "...": "..." } } ] }
Document: the passenger, the bookings inside it, each flight’s fields inside its booking
Read: a passenger’s trips in one document, one round trip, no matching
Write: a new booking added to this document, with the flight’s fields copied in
Copies: the departure time, in every passenger document with a booking on flight 7, and in flights
Schedule change: flight 7, 14:05 to 14:20
UPDATE, one rowHalf done: a failure after 90 writes
Facts to copy: read on every request, changed rarely
Facts to reference: changed while bookings exist
A copy saves a read on every request and adds a write on every change
{ "passenger_id": 42, "name": "A. Rivera", "bookings": [ { "booking_id": 9001, "seat": "14C", "fare_class": "Y", "flight_id": 7 }, { "booking_id": 9144, "seat": "2A", "fare_class": "J", "flight_id": 12 } ] }
Reference: flight_id in the booking, the flight’s fields in the flight document only
One copy: one write per schedule change, the new value on every read
Second read: the passenger document, then the flight documents it names
Embed what the read returns, reference what changes
Flight document with its bookings inside: every booking for flight 7 in one list
Document size limit: 16 MB, about 16,000 bookings at 1 KB each
Lists that grow: bookings per flight, position reports per aircraft, events per flight
flight_idLists that fit inside: read whole every time, with a small known bound
Data model: the shape given to the data
Model for a passenger’s trips: passenger document, bookings inside, flight fields copied in
Model for the seat sale: flight document with the count, booking documents of their own
Each model makes one access pattern a single round trip and the others several
Table per query: the same data stored in one shape per access pattern, each keyed for its query
Position reports: two questions, two tables
(aircraft_id, day), sorted by time, for “aircraft 3, last hour”(airport, minute), sorted by aircraft, for “every aircraft near LAX at 13:05”Copy writes
Cost: storage and writes, each multiplied by the number of copies
A question without its own copy is a scan
Document: one JSON object, stored, read and written as a unit
Entity: one document for one thing and what belongs to it
Identity: _id, one field every document has
Collection: a named set of documents
Instances
from pymongo import MongoClient db = MongoClient("mongodb://db.internal").airline # by a field, not a key db.passengers.find_one({"passenger_id": 42}) # {'_id': ObjectId('..'), 'passenger_id': 42, # 'name': 'A. Rivera', 'bookings': [...]} # a field inside an array db.passengers.find( {"bookings.flight.flight_id": 7}) # operator, and a projection db.flights.find( {"origin": "LAX", "departure_time": {"$gt": now}}, {"flight_number": 1, "departure_time": 1})
Filter: a document
{"$gt": now}, {"$in": [...]}
db.passengers.update_one( {"passenger_id": 42}, {"$push": {"bookings": { "booking_id": 9301, "flight_id": 7, "seat": "9F"}}}) db.flights.update_one( {"flight_id": 7, "seats_remaining": {"$gt": 0}}, # the check {"$inc": {"seats_remaining": -1}}) # the write # matched_count: 1 (sold) or 0 (full)
Update operators: $set a field, $inc a number, $push onto an array, $unset
Guarded update: one operation on one document
matched_count 0: no seat left, nothing writtenupdate_many: the same change to every match, atomic per document, not across documents
Atomic per document: a write to one document is all or nothing, and nothing partial is ever read
db.passengers.insert_one( {"passenger_id": 57, "name": "J. Okafor", "services": {"wheelchair": True}}) db.passengers.insert_one( {"passenger_id": 58, "name": "M. Chen"}) # no services field # absence handled in code doc.get("services", {}).get("wheelchair", False) db.command("collMod", "bookings", validator={"$jsonSchema": { "required": ["passenger_id", "flight_id", "fare_class"], "properties": { "fare_class": {"enum": ["F", "J", "W", "Y"]}, "flight_id": {"bsonType": "int"}}}})
Schemaless: no declared fields
Absent field: not an error, not NULL
Validator: a rule attached to the collection, checked on insert and update
flight_id 7 exists in flightsThe schema is in the application code, and in the validator where one is written
db.passengers.create_index( "passenger_id", unique=True) db.passengers.create_index( "bookings.flight.flight_id") # multikey db.flights.create_index( [("origin", 1), ("departure_time", 1)])
Without an index: find opens every document in the collection, as a table scan reads every row
Index on a field: one entry per document, as in a table
Multikey index: one entry per distinct array element
Compound index: origin, then departure time
Every index is written on every insert and on every update of its field
db.bookings.aggregate([ {"$match": {"passenger_id": 42}}, {"$lookup": { "from": "flights", "localField": "flight_id", "foreignField": "flight_id", "as": "flight"}}, {"$project": { "seat": 1, "fare_class": 1, "flight.origin": 1, "flight.departure_time": 1}}]) # one document per booking, # each with a flight array of one element
Pipeline: stages applied in order to a stream of documents
$match filters$lookup adds documents from another collection$project picks fields$group aggregatesCollections read: one by a filter, two by a lookup
Lookup: for each input document, the documents in from whose foreignField equals its localField, added as an array
flights.flight_id: one index lookup per bookingflights per bookingResult: the referenced documents read inside the server, returned in one round trip
Not a relational join: one collection added at a time, left outer by construction, no optimizer choosing the order
The N + 1 reads happen on the server, not across the network
with client.start_session() as s: with s.start_transaction(): r = db.flights.update_one( {"flight_id": 7, "seats_remaining": {"$gt": 0}}, {"$inc": {"seats_remaining": -1}}, session=s) if r.matched_count == 0: s.abort_transaction() else: db.bookings.insert_one( {"passenger_id": 42, "flight_id": 7, "seat": "9F"}, session=s) # both visible at commit, or neither
Multi-document atomicity: the all-or-nothing of one document, extended to several, on request
Transaction: operations on several documents and collections, committed together
Conflict: a second writer on flight 7 fails with a write conflict and retries
Default: no transaction
from pymongo import WriteConcern, ReadPreference sales = db.flights.with_options( write_concern=WriteConcern(w="majority")) board = db.flights.with_options( read_preference= ReadPreference.SECONDARY_PREFERRED) sales.update_one(...) # waits for 2 of 3 copies board.find({...}) # a copy a second behind
Replica set: one primary and its secondaries
Write concern: copies that acknowledge before the write returns
w=1: the primary onlyw="majority": a majority of the copiesRead preference: the copy that answers
One cluster serves the seat sale and the departures board at different levels, set per collection handle or per operation
Inside one document: fields found, indexed, changed together, validated
Referential integrity: none
flight_id: 7 in a booking names a flight documentEmbedded flight fields: no reference to another document
In a document store, the foreign-key check is application code or nothing
Table: typed columns, as in SQL; rows found by a two-part key, never by a join
Partition key: the first part; one value, one group of rows, one server
Clustering key: the second part; the order of the rows inside the group
Partition: one key’s rows, stored sorted; read whole or as a range of the clustering key
Wide row: the original form, thousands of columns per key; in CQL, many narrow rows sorted inside it
Implementations: Cassandra and ScyllaDB; Bigtable and HBase; Keyspaces on AWS, with the Cassandra interface
CREATE TABLE position_reports ( aircraft_id int, day date, reported_at timestamp, lat double, lon double, alt int, PRIMARY KEY ((aircraft_id, day), reported_at) ) WITH CLUSTERING ORDER BY (reported_at DESC); -- one aircraft, the last hour: one partition, one range SELECT reported_at, lat, lon, alt FROM position_reports WHERE aircraft_id = 3 AND day = '2026-10-06' AND reported_at > '2026-10-06 13:00:00'; -- the latest report: the first row of the partition SELECT * FROM position_reports WHERE aircraft_id = 3 AND day = '2026-10-06' LIMIT 1;
CQL: SQL’s shape, one table at a time
WHERE names the whole partition keyOne partition
Descending order: newest row first
Volume: a report every 5 s per aircraft
No partition key in the WHERE: the query is refused, not scanned
Placement: partition key hashed to a ring position
Whole: never split across servers
Growth: one key, one server, however many rows
Bucketing: a time unit in the partition key bounds the partition
The partition, fixed by its key before the first row, is the unit of placement, replication and reading
Write: appended to a log on disk and to a sorted table in memory, then acknowledged
Read: the partition’s pieces, from the memory table and from every file containing one, merged by clustering key
Delete: a tombstone appended
USING TTL on the insert: a tombstone with a date
No primary: every server accepts reads and writes
Replicas: the partition on N consecutive ring positions
Write: sent to all N, answered at the consistency level
ONE: one acknowledgementQUORUM: a majorityALL: every replicaRead: sent to the level’s number of replicas
Newest: by the timestamp each write carries
The consistency level is a keyword on every statement
One order: every row sorted by its full key, across the cluster
Tablets: contiguous key ranges
Range scan: adjacent keys on one server up to a tablet boundary
Row key design: what sorts together is read together
aircraft_id, day, reported_at: one aircraft’s history adjacentreported_at first: every aircraft’s current minute on one serverConsistency: one owner per range
DynamoDB, on AWS, also has one owner per key range
Row: one page
com.cnn.www and com.cnn.www/sports sort next to each otherColumn family: contents:, anchor:, language:, declared when the table is created; the unit of storage and of access control
Columns inside a family: created by writing them
A page linked from ten thousand others: ten thousand columns
A page linked from none: none
Versions: each cell stores its values by timestamp
contents:: the crawl historyanchor:: the input to link-based rankingRule: a table serves the queries that name a partition and read it in clustering order, and no others
Position reports: fit
(aircraft_id, day), sorted by time descending: one partition, one range, newest firstPassengers: do not fit
The table is written for its query before the first row arrives
Item: one document, JSON-shaped, up to 400 KB
{ "PK": "FLIGHT#7", "SK": "BOOKING#9001", "passenger_id": 42, "seat": "14C", "fare_class": "Y" }
Key
Operations: the whole interface
GetItem, PutItem, UpdateItem, DeleteItem: one item, by its full key; PutItem replaces the whole item, UpdateItem changes attributesQuery: one partition, in sort-key orderScan: every itemBatchWriteItem, TransactWriteItems: several itemsManaged service: regional
import boto3 from boto3.dynamodb.conditions import Key table = boto3.resource("dynamodb").Table("airline") # one partition, in SK order, up to 1 MB r = table.query( KeyConditionExpression= Key("PK").eq("FLIGHT#7") & Key("SK").begins_with("BOOKING#"), ReturnConsumedCapacity="TOTAL") r["Count"], r["ScannedCount"] # (12, 12) r["ConsumedCapacity"]["CapacityUnits"] # 1.0 # every item in the table, one page at a time r = table.scan( FilterExpression=Key("PK").eq("FLIGHT#7"), ReturnConsumedCapacity="TOTAL") r["Count"], r["ScannedCount"] # (12, 50000)
Query: a partition key value, plus a sort-key condition
LastEvaluatedKeyScan: every item in the table, page by page
Count: items returnedScannedCount: items readSame 12 bookings
Items a filter drops are read and billed
Read request unit
ConsistentRead=True per requestWrite request unit: 1 KB, rounded up per item
Billed quantity: bytes read or written, not items returned, not attributes named
Price: on-demand, $0.125 per million read units, $0.625 per million write units
The cost of a request is computable before it is sent
# the flight and its bookings: one partition table.put_item(Item={ "PK": "FLIGHT#7", "SK": "FLIGHT#7", "origin": "LAX", "departure_time": "2026-10-09T14:05", "seats_remaining": 1}) table.put_item(Item={ "PK": "FLIGHT#7", "SK": "BOOKING#9001", "passenger_id": 42, "seat": "14C", "fare_class": "Y"}) # one Query: the flight item and every booking table.query( KeyConditionExpression=Key("PK").eq("FLIGHT#7")) # a passenger's trips: a second copy, keyed the other way table.put_item(Item={ "PK": "PASSENGER#42", "SK": "BOOKING#9001", "flight_id": 7, "origin": "LAX", "departure_time": "2026-10-09T14:05", "seat": "14C"})
One table: flights, bookings, passengers as items side by side
Partition by the parent: FLIGHT#7 contains the flight item and its bookings
SK: the bookings aloneSort key as order
BOOKING#... sorts before FLIGHT#7SK orders reports by timeSecond copy: a passenger’s trips, keyed PASSENGER#42
Hot partition: a key every request shares, such as DATE#2026-10-06
PK, SK, GSI1PK, GSI1SK: the usual names for keys shared by several entity types
# the same items, keyed by passenger table.query( IndexName="ByPassenger", KeyConditionExpression= Key("GSI1PK").eq("PASSENGER#42") & Key("GSI1SK").between( "2026-10-01", "2026-12-31"))
GSI: a second table, maintained from the first by the service
Index write: a table write that touches the index’s keys is also a write to the index
Projection: the attributes the copy carries
# guarded update: condition and change, one item table.update_item( Key={"PK": "FLIGHT#7", "SK": "FLIGHT#7"}, UpdateExpression="SET seats_remaining = " "seats_remaining - :one", ConditionExpression="seats_remaining > :zero", ExpressionAttributeValues={":one": 1, ":zero": 0}) # ConditionalCheckFailedException: the seat is gone # the seat and the booking: two items, # together or not at all ddb = boto3.client("dynamodb") ddb.transact_write_items(TransactItems=[ {"Update": { "TableName": "airline", "Key": {"PK": {"S": "FLIGHT#7"}, "SK": {"S": "FLIGHT#7"}}, "UpdateExpression": "SET seats_remaining = seats_remaining - :one", "ConditionExpression": "seats_remaining > :zero", "ExpressionAttributeValues": {":one": {"N": "1"}, ":zero": {"N": "0"}}}}, {"Put": { "TableName": "airline", "Item": {"PK": {"S": "FLIGHT#7"}, "SK": {"S": "BOOKING#9301"}, "passenger_id": {"N": "42"}}}}])
Condition expression: checked and applied by the service as one operation on one item
Per item: put_item, update_item, delete_item atomic on their item
Transaction: up to 100 items across tables, all or none
Batch: 25 puts in one request
The seat and the booking: one transaction at twice the write units, or two writes with an interval where only one is applied
Per partition: at most 3,000 read units and 1,000 write units a second, whatever the table’s total
Throttled: past the rate, a request fails with ProvisionedThroughputExceededException
On-demand: billed per request unit
Provisioned: a per-second rate reserved and billed for the table
Split: a partition past its size or its rate divided in two by the service, later
A table’s capacity is usable only through keys that spread the load across partitions
Rule: the bill and the latency follow the bytes a request reads and the items a write touches
Departures board by Query: PK = AIRPORT#LAX, departure time in the sort key
Same board by Scan: every item, with a filter
No aggregation on the server: a sum or a count over the table is a Scan, read and billed in full
A table without a key for the access pattern is scanned, and billed for every item read
Node: an airport, with code, city, timezone as properties
Edge: a directed route from one airport to another, with minutes as a property
Path: a sequence of edges, such as LAX to DEN to JFK
Questions about paths
Connections: stored as edges
Instances
-- routes: one row per origin, destination pair served -- one stop: the table joined to itself SELECT r1.destination AS via FROM routes r1 JOIN routes r2 ON r2.origin = r1.destination WHERE r1.origin = 'LAX' AND r2.destination = 'JFK'; -- up to three edges: a recursive query WITH RECURSIVE reach AS ( SELECT destination AS airport, 1 AS stops FROM routes WHERE origin = 'LAX' UNION SELECT r.destination, reach.stops + 1 FROM reach JOIN routes r ON r.origin = reach.airport WHERE reach.stops < 3 ) SELECT airport, MIN(stops) AS stops FROM reach GROUP BY airport;
Routes as rows: a connection is a pair of codes
One stop: routes joined to itself
Any depth: a recursive query
Not in the result: the path itself
The rows give the destination, not the route
To avoid one airport on the way, or to sum the minutes, the recursion must carry the path along
The path is rebuilt from stored pairs on every read
// a node, an edge, a node: the pattern as drawn MATCH (lax:Airport {code: 'LAX'}) -[r:ROUTE]-> (jfk:Airport {code: 'JFK'}) RETURN r.minutes; // 320 // one stop: two edges, the stop in between MATCH (lax:Airport {code: 'LAX'}) -[:ROUTE]->(via:Airport) -[:ROUTE]->(jfk:Airport {code: 'JFK'}) RETURN via.code; // DEN // a path as one value MATCH p = (lax:Airport {code: 'LAX'}) -[:ROUTE]->()-[:ROUTE]-> (jfk:Airport {code: 'JFK'}) RETURN [n IN nodes(p) | n.code] AS stops, reduce(t = 0, e IN relationships(p) | t + e.minutes) AS minutes; // ["LAX", "DEN", "JFK"], 380
Pattern: (node)-[:EDGE]->(node)
Match: every place the pattern fits in the graph, one row per fit
Variables: a name on a node or edge, via, r, to return its properties
Path variable: p = ..., the whole path as a value
nodes(p), relationships(p): its parts, in orderThe store returns every path that fits the pattern
// everything within three edges of LAX MATCH (lax:Airport {code: 'LAX'}) -[:ROUTE*1..3]->(x:Airport) RETURN DISTINCT x.code; // SFO DEN DFW JFK (1 edge) // SEA ORD ATL (2 edges) // everything reachable at all MATCH (lax:Airport {code: 'LAX'})-[:ROUTE*]->(x) RETURN count(DISTINCT x); // the fewest stops to SEA MATCH p = shortestPath( (:Airport {code: 'LAX'}) -[:ROUTE*]->(:Airport {code: 'SEA'})) RETURN [n IN nodes(p) | n.code], length(p); // ["LAX", "SFO", "SEA"], 2
Expansion: one more edge from every node found, depth by depth
*1..3: stops at three edges*: stops at the end of the graphshortestPath: stops at the first arrival// a connection with a legal layover: a condition // between two consecutive edges MATCH (lax:Airport {code: 'LAX'}) -[f1:FLIGHT]->(via:Airport) -[f2:FLIGHT]->(jfk:Airport {code: 'JFK'}) WHERE f2.departs >= f1.arrives + duration({minutes: 45}) RETURN via.code, f1.number, f2.number, duration.between(f1.departs, f2.arrives) AS total ORDER BY total LIMIT 3; // any route of up to three edges that avoids ORD MATCH p = (:Airport {code: 'LAX'}) -[:ROUTE*1..3]->(:Airport {code: 'JFK'}) WHERE NONE(n IN nodes(p) WHERE n.code = 'ORD') RETURN [n IN nodes(p) | n.code]; // ["LAX", "JFK"] // ["LAX", "DEN", "JFK"] // ["LAX", "SFO", "DEN", "JFK"] // ["LAX", "DFW", "ATL", "JFK"]
Condition: between consecutive edges, or on every node of the path
Index-free adjacency
One hop: one read of the node’s edge list
Index lookup: once, for the starting node by code, then pointers only
Cost of a traversal: the edges touched, not the size of the graph
Cost of a join: an index lookup per row per hop
Joins are stored as pointers in a graph and recomputed on every read from a table
// the services of the airline application CREATE (web:Service {name: 'web'}), (api:Service {name: 'booking-api'}), (auth:Service {name: 'auth'}), (pay:Service {name: 'payments'}), (db:Service {name: 'postgres'}), (cache:Service {name: 'redis'}), (web)-[:DEPENDS_ON]->(api), (api)-[:DEPENDS_ON]->(auth), (api)-[:DEPENDS_ON]->(pay), (api)-[:DEPENDS_ON]->(db), (auth)-[:DEPENDS_ON]->(cache), (pay)-[:DEPENDS_ON]->(db); // everything that stops if auth stops MATCH (:Service {name: 'auth'}) <-[:DEPENDS_ON*1..]-(d:Service) RETURN DISTINCT d.name; // booking-api, web // the service most others depend on MATCH (s:Service)<-[:DEPENDS_ON*1..]-(d) RETURN s.name, count(DISTINCT d) AS dependents ORDER BY dependents DESC LIMIT 1; // postgres, 3
Reverse traversal: <-[:DEPENDS_ON*1..]-
Hub: a node with far more edges than the rest
postgres in a dependency graphPruning: checked as the pattern expands
*1..3WHERE r.minutes < 200Writes
Sharding: an edge across two servers makes every traversal over it a network round trip
The limit is the hub, not the size of the graph
Rule
LAX to JFK, two stops, not via ORD: the graph
Seats sold per route this month: the table
Both stores
The traversal is the one operation a graph store does better than a table
The booking page: one request, four reads
Three databases: four round trips
Three failure modes: any database down, the page down
Three to operate: a connection, a credential, a backup, a failover for each
The departures board: one airport’s flights in the next hours
In the relational database: a query on an index
In the key-value store: a cached copy per airport, 30 s lifetime
In DynamoDB: an item per flight under AIRPORT#LAX
The three differ in the board’s maximum age and in the cost of each write
-- a queue in the relational database UPDATE jobs SET state = 'running', worker = %s WHERE job_id = ( SELECT job_id FROM jobs WHERE state = 'queued' ORDER BY created_at FOR UPDATE SKIP LOCKED LIMIT 1) RETURNING job_id, payload; -- a document in the relational database SELECT services->>'wheelchair' FROM passengers WHERE passenger_id = 42; -- a path of fixed depth in the relational database SELECT r1.destination FROM routes r1 JOIN routes r2 ON r2.origin = r1.destination WHERE r1.origin = 'LAX' AND r2.destination = 'JFK';
A queue: one locked UPDATE* assigns a job to one worker
A document: the services object in a JSON column*
A path: two stops, two joins
Sessions: a table with an expiry column, up to the database’s lookup-rate limit
Second database: added when the first cannot serve the pattern at the rate, the size, or the depth needed
The database already running adds no connection, credential, backup, or failover
Original: the departure time in flights
Copies: the same value in board:LAX and in 180 trips items
No transaction across databases: the copy is written after the original, by a second write
Copy writes
Lag: the copy’s age at worst, plus whatever a failure leaves undone
Owner: the one database whose value of the fact is the truth
flightsCopy: the value in any other database
Path: how each copy is refreshed, and its lag
Rebuild: a wrong or lost copy is deleted and rebuilt from the owner
Two owners: the same fact written on its own in two databases
Rule
The booking: one database
The departures board: any of three
Unneeded additions
UPDATENo design makes every operation a single round trip