Object storage
Design a store for blobs by key at exabyte scale, where an acknowledged write survives disk, node and rack failures. The design keeps a strongly consistent index apart from the bytes.
- 1Split metadata from data: an index maps each key to a version and its fragment placement; storage nodes hold only fragments.
- 2Erasure coding 6+3 over racks, at most 2 fragments a rack: 1.5× the bytes, and a rack plus a node can fail.
- 3Durability is a race between disk failures and repair, so repair bandwidth and scrubbing are part of the design.
- 4One latest row per key is the commit point: strong read-after-write, prefix listing as a range scan, and safe GC.
Design an object store: put and get blobs by key, eleven nines of durability, at exabyte scale.
- functional
- Buckets; put, get, head, delete, list by prefix; versions; multipart upload.
- sizes
- Objects from 0 B to 5 TB; 1 MB average.
- durable
- 11 nines a year per object.
- available
- 99.99% for reads; a rack can fail.
- consistency
- Strong read-after-write and list-after-write.
- scale
- 1 EB stored, 2 PB more each day.
I will keep a strongly consistent index of keys apart from the bytes, and erasure-code the bytes across racks.
| question | answer assumed | it decides |
|---|---|---|
| How large is an average object? | 1 MB; many under 64 KB | Object count, index size |
| Reads and writes a day? | 2 × 10⁹ PUTs, 10 GETs per PUT | Front ends, index shards |
| Must a list see a write at once? | Yes: strong list-after-write | One index, range shards |
| Overwrite or append? | Whole-object overwrite only | Immutable versions, GC |
| What must fail without loss? | A rack and one more disk | Code, placement |
| Keep old versions? | Per bucket: on or off | GC rule |
| One region or many? | One region; copy to a second on request | Replication scope |
I ask about object sizes, the read-write mix, listing consistency and the failures to survive. Each answer changes the layout.
- objects
- 1 EB / 1 MB = 10¹² objects.
- PUTs
- 2 PB a day / 1 MB / 86,400 s ≈ 23,148 a second; 3× peak ≈ 69,444.
- GETs
- 10 × PUTs ≈ 231,481 a second; 3× peak ≈ 694,444.
- bandwidth
- In: 23,148 × 1 MB ≈ 23 GB/s. Out: ≈ 231 GB/s average.
- disks
- 1 EB × 1.5 (6+3) / 20 TB = 75,000 disks.
- failures
- 2% a year × 75,000 = 1,500 a year ≈ 4.1 a day.
- repair
- Each failed disk: 6 × 20 TB = 120 TB read. A day: ≈ 493 TB ≈ 5.7 GB/s, all the time.
- index
- 10¹² × 500 B (key, version, 9 placements) = 500 TB; at 2 TB a shard, 250 shards.
Storage decides the index shard count, not throughput: one lab shard commits about 5,000 PUTs a second.
About 75,000 disks at 1.5× overhead, 4 disk failures a day, and repair reads of about 493 TB a day.
| layout | bytes | survives | 2%, 24 h | 5%, 72 h |
|---|---|---|---|---|
| 3 replicas | 3× | 2 | 10.2 | 8.1 |
| Reed-Solomon 6+3 | 1.5× | 3 | 12.4 | 9.4 |
| Reed-Solomon 10+4 | 1.4× | 4 | 15.4 | 11.5 |
- model
- p = AFR × repair hours / 8,760. An object is lost when more than m of its k + m disks fail in one window.
- 11 nines
- A yearly loss chance of 10⁻¹¹ per object. × 10¹² objects = 10 objects a year.
- 6+3
- 4.1 × 10⁻¹³ × 10¹² ≈ 0.4 objects a year.
Columns: disk failures a year, rebuild time. Nines are derived; the model ignores correlated failures, which placement must remove.
6+3 with a 24-hour rebuild gives 12.4 nines in the independent-failure model. Placement over racks is what makes failures independent enough for the model to hold.
| call | success | errors |
|---|---|---|
| PUT /{bucket}/{key} + Content-SHA256 | 200 ETag, version-id | 400 bad digest, 403, 503 |
| GET /{bucket}/{key}[?version-id] + Range | 200 or 206 bytes | 404, 503 |
| HEAD /{bucket}/{key} | 200 size, ETag | 404 |
| DELETE /{bucket}/{key} | 204, delete marker | 403 |
| GET /{bucket}?prefix=&after=&limit= | 200 keys + next token | 403 |
| POST /{bucket}/{key}?uploads | 200 upload-id | 403 |
| PUT ...?upload-id=&part=n | 200 part ETag | 400 bad digest |
| POST ...?upload-id= + part list | 200 ETag, version-id | 400 missing or small part |
- A retried PUT writes the same bytes as a new version. The result is the same object, so a retry is safe.
- DELETE is idempotent: a second delete adds a marker over a marker.
- The next token is the last key of the page. New keys do not shift later pages, so no existing key is skipped or repeated.
id = create_upload(bucket, key) // nothing visible yet1
PARALLEL FOR EACH part n: // 5 MiB to 5 GiB each
upload_part(id, n, bytes, sha256) // retry one part alone2
complete(id, [1, 2, ..., N]) // one commit: all or nothing3- 1Parts live apart from the key. A GET still returns the previous version.
- 2A network failure costs one part, not the whole upload. A part with a bad checksum is refused.
- 3Complete joins the listed parts and commits one version. Abandoned uploads expire through a lifecycle rule.
Tested source Go: multipart, reused from the B9 lab
// CreateUpload starts a multipart upload for key. Nothing is visible under key until Complete.
func (s *ObjectStore) CreateUpload(key string) string {
s.mu.Lock()
defer s.mu.Unlock()
s.nextID++
id := "u" + strconv.Itoa(s.nextID)
s.uploads[id] = &upload{key: key, parts: map[int][]byte{}}
return id
}
// UploadPart stores one part. Parts can arrive in any order and in parallel; sending a part again
// replaces it. The client sends the SHA-256 of the part, and a corrupted part is refused.
func (s *ObjectStore) UploadPart(id string, n int, data []byte, sum string) error {
if digest(data) != sum {
return fmt.Errorf("part %d: %w", n, ErrBadChecksum)
}
s.mu.Lock()
defer s.mu.Unlock()
u, ok := s.uploads[id]
if !ok {
return fmt.Errorf("upload %s: %w", id, ErrNoUpload)
}
u.parts[n] = slices.Clone(data)
return nil
}
// ListParts returns the part numbers the store holds, so a client can resume after a crash.
func (s *ObjectStore) ListParts(id string) ([]int, error) {
s.mu.Lock()
defer s.mu.Unlock()
u, ok := s.uploads[id]
if !ok {
return nil, fmt.Errorf("upload %s: %w", id, ErrNoUpload)
}
nums := make([]int, 0, len(u.parts))
for n := range u.parts {
nums = append(nums, n)
}
slices.Sort(nums)
return nums, nil
}
// Complete joins the listed parts in order into one object. The object appears all at once.
func (s *ObjectStore) Complete(id string, parts []int) (Object, error) {
s.mu.Lock()
defer s.mu.Unlock()
u, ok := s.uploads[id]
if !ok {
return Object{}, fmt.Errorf("upload %s: %w", id, ErrNoUpload)
}
var data []byte
for i, n := range parts {
p, ok := u.parts[n]
if !ok {
return Object{}, fmt.Errorf("part %d: %w", n, ErrMissingPart)
}
if i < len(parts)-1 && len(p) < s.MinPart {
return Object{}, fmt.Errorf("part %d is %d bytes: %w", n, len(p), ErrSmallPart)
}
data = append(data, p...)
}
o := Object{Data: data, SHA256: digest(data), ContentType: "application/octet-stream"}
s.objects[u.key] = o
delete(s.uploads, id)
return o, nil
}
A PUT replaces the whole object and returns a version id. Large objects go up in parts that commit together, so a reader never sees half an object.
- versions
- Immutable. An overwrite adds a row; it never edits one.
- latest
- The only row a commit changes. A read and a list both go through it.
- nodes
- Keep fragment bytes by version and index. They know nothing about keys.
CREATE TABLE buckets (
id bigserial PRIMARY KEY,
name text NOT NULL UNIQUE,
versioned boolean NOT NULL DEFAULT false
);
-- Each PUT or DELETE adds a version; it is written before it is visible.
CREATE TABLE versions (
id bigserial PRIMARY KEY,
bucket_id bigint NOT NULL REFERENCES buckets,
key text COLLATE "C" NOT NULL1,
size_bytes bigint NOT NULL,
etag text NOT NULL,
is_delete boolean NOT NULL DEFAULT false,
created_at timestamptz NOT NULL DEFAULT now(),
committed boolean NOT NULL DEFAULT false2
);
CREATE INDEX versions_key ON versions (bucket_id, key, id DESC);
-- Where each erasure-coded fragment of a version lives.
CREATE TABLE fragments (
version_id bigint NOT NULL REFERENCES versions ON DELETE CASCADE,
idx smallint NOT NULL,
node_id text NOT NULL,
checksum bytea NOT NULL3,
PRIMARY KEY (version_id, idx)
);
-- One row per key: the commit point. Reads see the version it names.
CREATE TABLE latest (
bucket_id bigint NOT NULL REFERENCES buckets,
key text COLLATE "C" NOT NULL,
version_id bigint NOT NULL UNIQUE REFERENCES versions4,
PRIMARY KEY (bucket_id, key)
);- 1Byte order, not language order, so a prefix is one contiguous range of the index.
- 2A version is written before it is visible. A versioned bucket keeps committed versions; GC takes the rest.
- 3Scrub compares each fragment with this. It also catches a disk that returns the wrong block.
- 4The foreign key stops GC from deleting a version a key points to.
The index has one row per version, one row per fragment, and one latest row per key. A commit is one upsert of that latest row.
Step 1: Write
- The front end checks auth and the checksum, then splits the bytes into 6 data and 3 parity fragments.
- It writes each fragment to a node chosen by the placement rule: at most 2 in a rack.
- After 9 acks it adds the version rows, then commits: one upsert of the key's latest row.
If it fails
A node does not ack: write that fragment to another node. The front end dies before commit: the version stays invisible, and GC deletes it after the grace period.
The front end encodes and writes 9 fragments, then commits one index row. Workers repair, scrub and collect garbage in the background.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Front end | Reed-Solomon k + m | Any k of k + m fragments rebuild the object; 1.5× bytes for 6+3. | RAID 6, backups, QR codes |
| Placement | At most ⌈(k+m)/racks⌉ fragments a rack | A rack failure takes at most 2 of 9. | Replica placement in any database |
| Storage node | A checksum per fragment | A bad fragment is skipped on read and found by scrub. | File systems, network frames |
| Workers | Repair from any k | Spare fragments come back before the next failure. | Distributed file systems |
| Postgres | INSERT ... ON CONFLICT ... WHERE | The commit point: the newest version wins, read-after-write holds. | Optimistic writes, idempotency |
| Postgres | B-tree in byte order | A prefix list is one index range scan. | Time ranges, autocomplete |
| Postgres | Foreign keys | GC cannot delete a live version. | Any referential rule |
| Postgres | Range shards with sync replicas | The index grows past one machine and survives a primary loss. | Large key-value indexes |
| CDN | Edge cache | Hot public objects are served without touching the nodes. | Static sites, video |
Erasure coding and placement give durability, checksums and repair keep it, and one Postgres row per key gives consistency.
| placement | status | rack 1 fails |
|---|---|---|
| 6+3 on 9 random nodes | Not approved | 11 of 1,000 lost |
| 3 replicas on 3 random nodes | Not approved | 3 of 1,000 lost |
| 6+3, at most 2 a rack | Approved | 0 lost; also with 1 more node |
| 6+3, 3 a zone over 3 zones | A zone must fail safely | Every write and repair crosses zones |
Lab: 6 racks of 4 nodes, 1,000 objects of 6 KiB. A test fails each rack together with each other node in turn.
put(bucket, key, bytes):
fragments = encode(bytes, k = 6, m = 3) // 9 pieces of size / 61
FOR EACH fragment i:
node = place(i) // at most 2 per rack2
write fragment i to node, wait for ack
id = index.begin(bucket, key, nodes, checksums)
index.commit(id) // now visible3
RETURN 200, version id
get(bucket, key):
v = index.latest(bucket, key) // one row, on the primary4
read fragments 0..5 of v
FOR EACH one that is down or fails its checksum:
read the next parity fragment instead5
IF fewer than 6 good: RETURN 5036
RETURN decode(the 6)- 1The 6 data fragments are the object cut in 6. A read with all of them up needs no decoding.
- 2The rule prefers the rack with the fewest fragments, and skips down nodes.
- 3Fragments first, commit last. A crash in between leaves garbage, never a broken object.
- 4Reads of the index go to the primary of its shard. A replica could return the version before the commit.
- 5A degraded read. In the lab, one node down made 267 of 1,000 reads degraded and none failed.
- 6With more than m fragments gone the read fails. It never returns wrong bytes.
Tested source Go: place · Go: put and commit · Go: get · Go: Reed-Solomon rebuild, reused from the B9 lab
// place picks a node for one more fragment of v. Spread only takes a node whose domain holds
// fewer than Cap fragments of v, and prefers the domain that holds the fewest. Both policies
// skip down nodes and nodes that already hold a fragment of v.
func (c *Cluster) place(v *Version) (int, error) {
used := map[int]bool{}
inDomain := make([]int, c.Domains)
for _, f := range v.Frags {
if f.Node >= 0 {
used[f.Node] = true
inDomain[c.Nodes[f.Node].Domain]++
}
}
var cands []int
best := c.Code.K + c.Code.M + 1
for i, n := range c.Nodes {
if !n.Up || used[i] {
continue
}
if c.Policy == Spread {
d := inDomain[n.Domain]
if d >= c.Cap() || d > best {
continue
}
if d < best { // a domain with fewer fragments: start the list again
best, cands = d, cands[:0]
}
}
cands = append(cands, i)
}
if len(cands) == 0 {
return -1, fmt.Errorf("fragment of %s: %w", v.Key, ErrNoPlacement)
}
return cands[c.rng.IntN(len(cands))], nil
}
// Put encodes data into K+M fragments and writes each to its own node. The version is not
// visible until Commit: a read in between still returns the previous version.
func (c *Cluster) Put(key string, data []byte) (*Version, error) {
shards := c.Code.Encode(data)
c.nextID++
v := &Version{Key: key, ID: c.nextID, Size: len(data), FragSize: len(shards[0]), Frags: make([]Frag, len(shards))}
for i := range v.Frags {
v.Frags[i].Node = -1
}
for i, s := range shards {
n, err := c.place(v)
if err != nil {
c.drop(v) // fragments already written are garbage now
return nil, err
}
c.Nodes[n].frags[fragID{v.ID, i}] = s
v.Frags[i] = Frag{Node: n, Sum: sha256.Sum256(s)}
}
c.versions[v.ID] = v
return v, nil
}
// Commit makes v the latest version of its key: one write to the index. The version it
// replaces stays on the nodes until GC.
func (c *Cluster) Commit(v *Version) {
v.Committed = true
c.latest[v.Key] = v
}
// Get reads the latest version of key. It reads the K data fragments; for each one it cannot
// read, it reads a parity fragment instead and decodes. With more than M fragments gone it
// returns ErrUnreadable, never wrong bytes.
func (c *Cluster) Get(key string) ([]byte, error) {
v, ok := c.latest[key]
if !ok || v.Deleted {
return nil, fmt.Errorf("%s: %w", key, ErrNotFound)
}
shards, degraded, err := c.readK(v)
if err != nil {
return nil, err
}
c.Traffic.Reads++
c.Traffic.ReadBytes += int64(c.Code.K * v.FragSize)
if degraded {
c.Traffic.DegradedReads++
}
return slices.Concat(shards[:c.Code.K]...)[:v.Size], nil
}
// readK collects K good fragments, data fragments first, and rebuilds the rest.
func (c *Cluster) readK(v *Version) (shards [][]byte, degraded bool, err error) {
shards = make([][]byte, len(v.Frags))
have := 0
for i := range v.Frags {
if have == c.Code.K {
break
}
if b, ok := c.fragment(v, i); ok {
shards[i] = slices.Clone(b)
have++
} else if i < c.Code.K {
degraded = true
}
}
if have < c.Code.K {
return nil, false, fmt.Errorf("%s: %d of %d fragments: %w", v.Key, have, len(v.Frags), ErrUnreadable)
}
if degraded {
if err := c.Code.Reconstruct(shards); err != nil {
return nil, false, fmt.Errorf("decode %s: %w", v.Key, err)
}
}
return shards, degraded, nil
}
// Reconstruct fills in the missing (nil) shards. Any K surviving shards are K known rows of the
// encoding matrix times the data, so inverting those K rows recovers the data shards. The lost
// parity shards are then encoded again.
func (c *Code) Reconstruct(shards [][]byte) error {
var rows matrix
var have [][]byte
for i, s := range shards {
if s != nil && len(rows) < c.K {
rows = append(rows, c.enc[i])
have = append(have, s)
}
}
if len(rows) < c.K {
return fmt.Errorf("%d of %d shards left, need %d: %w", len(rows), c.K+c.M, c.K, ErrTooFewShards)
}
dec, err := rows.invert()
if err != nil {
return fmt.Errorf("decode: %w", err)
}
data := make([][]byte, c.K)
for j := range c.K {
data[j] = c.combine(dec[j], have)
}
for i := range shards {
if shards[i] == nil {
shards[i] = c.combine(c.enc[i], data)
}
}
return nil
}
I spread the 9 fragments over 6 racks, at most 2 in each. In the lab, random placement lost 11 of 1,000 objects to one rack; spread placement lost none.
| method | status | cost |
|---|---|---|
| Rebuild only when an object is read | Not approved | Cold objects lose spares in silence. |
| Rebuild every object on a dead node, from any 6 | Approved | 3.3 B read per B rebuilt |
| 3 replicas: copy a survivor | Hot or small data | 0.9 B read per B rebuilt |
| Wait until 2 of 3 spares are gone | Repair bandwidth is short | Less traffic; one spare left at times |
| Scrub every disk every few weeks | Approved | Reads every byte; finds bit rot |
- Repair reads from many nodes and writes to many nodes, so a dead disk rebuilds in parallel, not from one neighbour.
- Repair objects with the fewest spare fragments first.
- Mark a node dead after a timeout of minutes, so a reboot does not start a full rebuild.
Lab: after a rack and a node failed, repair fixed all 1,000 objects: 1,838 fragments, 6,000 KiB read for 1,838 KiB written.
repair(): // workers, all the time
FOR EACH version with a fragment on a dead node:
IF more than m fragments are gone: report lost1
read any k good fragments // k reads per object2
rebuild the missing ones
write each to a new node that place() allows3
update the version's fragment rows
scrub(): // every disk, every few weeks4
FOR EACH fragment on the disk:
IF checksum(bytes) != stored checksum:
delete it // repair rebuilds it5- 1This must stay at zero. A lost object is a page to the on-call engineer, not a log line.
- 2One read of 6 fragments rebuilds all the missing ones of that object.
- 3The same rule as a write, so the repaired object again survives a rack.
- 4A disk can return bad bytes with no error. Scrub finds them while spares still exist.
- 5In the lab, scrub found 10 corrupt fragments of 10, and repair rebuilt all 10.
Tested source Go: repair · Go: scrub · Go: GC on the nodes · Go: durability model, reused from the B9 lab
// Repair rebuilds every missing fragment of every committed version. For one version it reads
// K good fragments once, decodes, and writes each missing fragment to a new node that the
// placement rule allows. The old copy, if its node returns, is an orphan for GC.
func (c *Cluster) Repair() (RepairResult, error) {
var res RepairResult
for _, id := range c.ids() {
v := c.versions[id]
if !v.Committed || v.Deleted {
continue
}
var gone []int
for i := range v.Frags {
if _, ok := c.fragment(v, i); !ok {
gone = append(gone, i)
}
}
if len(gone) == 0 {
continue
}
res.Objects++
if len(gone) > c.Code.M {
res.Lost++
continue
}
shards, _, err := c.readK(v)
if err != nil {
return res, err
}
if err := c.Code.Reconstruct(shards); err != nil { // fills the parity readK skipped
return res, fmt.Errorf("rebuild %s: %w", v.Key, err)
}
c.Traffic.RepairRead += int64(c.Code.K * v.FragSize)
for _, i := range gone {
v.Frags[i].Node = -1 // the old nodes no longer count in the placement rule
}
for _, i := range gone {
n, err := c.place(v)
if err != nil {
return res, err
}
c.Nodes[n].frags[fragID{v.ID, i}] = shards[i]
v.Frags[i].Node = n
c.Traffic.RepairWrite += int64(v.FragSize)
c.Traffic.Rebuilt++
res.Rebuilt++
}
}
return res, nil
}
// Scrub reads every fragment on every healthy node and checks it against the checksum in the
// index. A fragment that fails is deleted, so the next repair rebuilds it from the others.
func (c *Cluster) Scrub() int {
bad := 0
for _, id := range c.ids() {
v := c.versions[id]
for i, f := range v.Frags {
if f.Node < 0 || !c.Nodes[f.Node].Up {
continue
}
key := fragID{v.ID, i}
if b, ok := c.Nodes[f.Node].frags[key]; ok && sha256.Sum256(b) != f.Sum {
delete(c.Nodes[f.Node].frags, key)
bad++
}
}
}
return bad
}
// GC deletes what no reader can reach: every committed version that is not the latest of its
// key (overwritten, or hidden by a delete marker), and every fragment a node holds that the
// index places elsewhere (left behind when repair moved it). It returns the bytes freed.
func (c *Cluster) GC() int64 {
var freed int64
for _, id := range c.ids() {
v := c.versions[id]
if v.Committed && c.latest[v.Key] != v {
for i, f := range v.Frags {
if f.Node >= 0 {
freed += int64(len(c.Nodes[f.Node].frags[fragID{v.ID, i}]))
}
}
c.drop(v)
}
}
for ni, n := range c.Nodes {
if !n.Up {
continue // swept when the node returns
}
for k, b := range n.frags {
v, ok := c.versions[k.version]
if !ok || v.Frags[k.idx].Node != ni {
freed += int64(len(b))
delete(n.frags, k)
}
}
}
return freed
}
// Durable works out the yearly chance to lose an object. The model: disks fail independently at
// rate afr a year, and a failed disk is rebuilt within repairHours. An object is lost only if more
// than M of its K+M disks fail inside one repair window.
func Durable(s Scheme, afr, repairHours float64) Durability {
n := s.K + s.M
p := afr * repairHours / 8760 // one disk, one window
var perWindow float64
for i := s.M + 1; i <= n; i++ {
perWindow += binom(n, i) * math.Pow(p, float64(i)) * math.Pow(1-p, float64(n-i))
}
windows := 8760 / repairHours
annual := windows * perWindow
return Durability{
Overhead: float64(n) / float64(s.K),
Tolerates: s.M,
P: p,
Windows: windows,
Lead: binom(n, s.M+1) * math.Pow(p, float64(s.M+1)),
PerWindow: perWindow,
Annual: annual,
Nines: -math.Log10(annual),
}
}
Repair reads 6 fragments to rebuild any of them, so 6+3 moves about 3.3 bytes per byte it rebuilds. At this scale that is several GB/s of repair traffic, all the time.
| index | status | prefix list |
|---|---|---|
| One Postgres | Up to a few TB of rows | One range scan |
| Hash-sharded key-value store | Not approved | Asks every shard |
| Range shards by (bucket, key) | Approved | One shard or neighbours |
| Consensus-replicated key ranges | Many regions | One range, quorum writes |
| Eventually consistent replicas | Not approved | Can miss a new key |
- Keys that start with a date or a counter put all writes on the last range. Split that range, or add a hash before the key.
- In the lab, 16 writers committed one key in shuffled order; the newest version won every time.
- A pending upload inside the grace period was never collected.
One lab shard: about 5,000 commits, 80,000 HEADs and 15,000 lists of 100 keys a second, 32 clients.
commit(id):
UPSERT latest(bucket, key) = id
only IF the stored id < id1 // the newest write wins
list(bucket, prefix, after, limit):
scan latest in key order
FROM max(prefix, after) TO prefix_end2
skip delete markers; stop at limit
RETURN keys, next token = last key3
gc(): // a job, every hour
FOR EACH version no latest row names,
older than the grace period4:
delete its fragments on the nodes
delete its rows // FK refuses a live one5- 1Two writers commit in any order; the larger version id stays. The other version is garbage.
- 2prefix_end is the prefix with its last byte raised by one. The scan touches only matching keys.
- 3A key token, not an offset, so new keys do not shift the pages.
- 4Longer than the slowest write, so GC never takes a version between write and commit.
- 5A bug in GC then fails loudly instead of deleting a live object.
Tested source SQL: commit · SQL: list · SQL: garbage · Go: begin and commit · Go: list · Go: garbage
-- The newest version wins, whatever order two writers commit in.
WITH v AS (
UPDATE versions SET committed = true WHERE id = $1 RETURNING bucket_id, key, id
)
INSERT INTO latest (bucket_id, key, version_id)
SELECT bucket_id, key, id FROM v
ON CONFLICT (bucket_id, key) DO UPDATE SET version_id = EXCLUDED.version_id
WHERE latest.version_id < EXCLUDED.version_id
RETURNING version_id;-- One page of keys under a prefix, in key order, after the last key of the previous page.
SELECT l.key, v.size_bytes, v.etag
FROM latest l JOIN versions v ON v.id = l.version_id
WHERE l.bucket_id = $1
AND l.key >= $2 AND ($3::text IS NULL OR l.key < $3)
AND l.key > $4
AND NOT v.is_delete
ORDER BY l.key
LIMIT $5;-- Versions no reader can reach, older than the grace period: overwritten, hidden by a delete
-- marker, or written and never committed. A versioned bucket keeps its old versions.
SELECT v.id
FROM versions v JOIN buckets b ON b.id = v.bucket_id
WHERE NOT EXISTS (SELECT 1 FROM latest l WHERE l.version_id = v.id)
AND v.created_at < now() - make_interval(secs => $1)
AND (NOT b.versioned OR NOT v.committed)
ORDER BY v.id
LIMIT $2;// Begin records a version and where its fragments went, in one transaction, after the data
// nodes acknowledged them. Nothing reads it until Commit.
func (m Meta) Begin(ctx context.Context, bucket int64, key string, size int64, etag string, frags []Placed) (int64, error) {
var id int64
err := pgx.BeginFunc(ctx, m.DB, func(tx pgx.Tx) error {
if err := tx.QueryRow(ctx, stmts["begin_version"], bucket, key, size, etag, false).Scan(&id); err != nil {
return fmt.Errorf("add version: %w", err)
}
for i, f := range frags {
if _, err := tx.Exec(ctx, stmts["add_fragment"], id, i, f.Node, f.Checksum); err != nil {
return fmt.Errorf("add fragment %d: %w", i, err)
}
}
return nil
})
if err != nil {
return 0, fmt.Errorf("begin %s: %w", key, err)
}
return id, nil
}
// Commit points the key at version id, unless a newer version already holds it. It returns
// false when a newer write won; the caller's version is then garbage.
func (m Meta) Commit(ctx context.Context, id int64) (bool, error) {
var got int64
err := m.DB.QueryRow(ctx, stmts["commit"], id).Scan(&got)
if errors.Is(err, pgx.ErrNoRows) {
return false, nil
}
if err != nil {
return false, fmt.Errorf("commit version %d: %w", id, err)
}
return true, nil
}
// List returns up to limit keys that start with prefix and sort after the key `after`. The next
// page passes the last key of this one. Keys sort by byte value, so one prefix is one range of
// the primary key index.
func (m Meta) List(ctx context.Context, bucket int64, prefix, after string, limit int) ([]Entry, error) {
var upper *string
if u, ok := prefixEnd(prefix); ok {
upper = &u
}
rows, err := m.DB.Query(ctx, stmts["list"], bucket, prefix, upper, after, limit)
if err != nil {
return nil, fmt.Errorf("list %q: %w", prefix, err)
}
var out []Entry
for rows.Next() {
var e Entry
var etag string
if err := rows.Scan(&e.Key, &e.Size, &etag); err != nil {
return nil, fmt.Errorf("list %q: %w", prefix, err)
}
out = append(out, e)
}
return out, rows.Err()
}
// prefixEnd returns the smallest string above every string that starts with prefix: the
// prefix with its last byte raised by one. An empty prefix has no end.
func prefixEnd(prefix string) (string, bool) {
b := []byte(prefix)
for i := len(b) - 1; i >= 0; i-- {
if b[i] < 0x7f {
b[i]++
return string(b[:i+1]), true
}
}
return "", false
}
// Garbage returns versions that no reader can reach and that are older than grace. The grace
// period protects a write between Begin and Commit.
func (m Meta) Garbage(ctx context.Context, grace time.Duration, limit int) ([]int64, error) {
rows, err := m.DB.Query(ctx, stmts["garbage"], grace.Seconds(), limit)
if err != nil {
return nil, fmt.Errorf("find garbage: %w", err)
}
ids, err := pgx.CollectRows(rows, pgx.RowTo[int64])
if err != nil {
return nil, fmt.Errorf("find garbage: %w", err)
}
return ids, nil
}
// Purge deletes a version's rows after its fragments are gone from the nodes. Postgres refuses
// if a key still points to it.
func (m Meta) Purge(ctx context.Context, id int64) error {
if _, err := m.DB.Exec(ctx, stmts["purge"], id); err != nil {
return fmt.Errorf("purge version %d: %w", id, err)
}
return nil
}
I keep the index in range-sharded Postgres with synchronous replicas. A commit is one upsert, so a read or a list right after a write sees it.
1,000 objects of 6 KiB · 9 fragments each of 1 KiB · any 6 rebuild the object · at most 2 per rack · 1.5× bytes stored
Write 1,000 objects. Every fragment is on a healthy node.
bars: fragments each node holdsf0 to f8: object photos/2026/0172.jpg, 0 of 9 fragments gone, readable
Spread placement keeps every object readable through a rack and a node. Random placement lost 47 of 1,000 objects to the same failures.
| event | result | why it is safe | saved by |
|---|---|---|---|
| A disk or node dies | Reads of its fragments become degraded. | Any 6 of 9 fragments rebuild the object. Repair restores the spares. | Erasure code |
| A rack loses power | Up to 2 fragments of each object are gone. | 7 remain, one more than needed. | Placement |
| A disk returns bad bytes | One fragment fails its checksum. | The read skips it; scrub deletes it; repair rebuilds it. | Checksums |
| The front end dies before commit | 9 fragments and an uncommitted version. | No latest row names it. GC deletes it after the grace period. | Commit point |
| Two clients write one key at once | Two versions. | The commit keeps the larger version id; the other becomes garbage. | Conditional upsert |
| An index primary fails | Writes to its key range pause. | A synchronous replica has every commit and takes over in seconds. | Sync replica |
| A batch of disks fails together | Failures are no longer independent. | Mix disk models and ages across racks; the model assumes independence. | Operations |
| One object gets most of the reads | Its 6 data nodes run hot. | Serve it from a CDN, or read different sets of 6 of its 9 fragments. | CDN |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One server: files on local disks, the index in Postgres. | Tens of TB; about 5,000 PUTs a second at the index. | One disk or server failure must not lose data or stop reads. |
| 2 | 3 replicas on 3 servers; a sync standby for the index. | Hundreds of TB; survives 2 lost copies. | Storage cost: every byte is stored 3 times. |
| 3 | Erasure coding 6+3 over 6 racks, with repair and scrub. | Petabytes at 1.5×; a rack and a node can fail. | The index outgrows one Postgres. |
| 4 | Range-sharded index, split by (bucket, key). | About 250 shards for 10¹² objects. | One cluster's repair traffic, network or racks hit a limit. |
| 5 | Cells: many clusters; a directory maps each bucket to a cell. | Exabytes; a cell failure touches only its buckets. | Top of the ladder; add a second region per bucket on request. |
Shard capacity is derived: the lab's one-shard rate × 250 shards. Demand is from panel C.
I start with files on one server and Postgres for the index. I move to erasure coding when copies cost too much, and to cells when one cluster cannot grow.
| question | answer | go deeper |
|---|---|---|
| How do you store billions of tiny objects? | Pack many small objects into one large erasure-coded file, and keep the offset in the index. One fragment set then serves thousands of objects. | B9 |
| How do clients upload without the front end in the path? | The service signs a URL for one key and a short time. The client sends bytes straight to the store. | B9 |
| One object is read a million times a second. Now what? | Put a CDN in front for public objects. Inside the store, add replicas for hot objects or read different sets of 6 fragments. | B5 |
| Why is read-after-write strong here? | Every read and every commit go to the one latest row on its shard primary. There is no replica to read stale data from. | X3 |
| How do you split the index? | By key range, with automatic splits of hot or large ranges. A hash prefix spreads keys that start with a date. | X2 |
| How do you copy a bucket to another region? | Ship committed versions in order, after the commit. The copy lags by seconds, and a region loss loses only that lag. | X1 |
0 of 9 known
Why Reed-Solomon 6+3 and not 3 replicas?
There are 9 fragments and 6 racks. Why at most 2 in a rack?
A front end wrote all 9 fragments, then crashed before the commit. What do readers see?
Why does repair of a 6+3 object move more bytes than repair of a replica?
Eleven nines, 10¹² objects. How many objects are lost a year?
Why does repair speed decide durability?
Why shard the index by key range and not by hash?
How does GC avoid deleting an object a reader can see?
What does scrub find that repair does not?
- nines
- 6+3 at 2% AFR, 24 h rebuild: 12.4. 3 replicas: 10.2.
- bytes
- 6+3: 1.5×. 10+4: 1.4×. 3 replicas: 3×.
- 11 nines
- 10¹² objects lose about 10 a year.
- placement
- One rack down: random lost 11 of 1,000; spread lost 0.
- repair
- 6+3 reads 3.3 B per B rebuilt; replicas 0.9. A 20 TB disk: 120 TB read.
- degraded
- One node of 24 down: 27% of reads used parity, 0 failed.
- index
- One Postgres: about 5,000 commits and 80,000 HEADs a second.
- S3
- Designed for 11 nines; strong read-after-write.
Cluster runs are seeded and repeat exactly. Index rates: Postgres 16 on a 16-thread laptop shared with other jobs, 32 clients. S3 figures are cited from its documentation.