System Design
C11

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.

Not startedSaved in this browser only.
  1. 1Split metadata from data: an index maps each key to a version and its fragment placement; storage nodes hold only fragments.
  2. 2Erasure coding 6+3 over racks, at most 2 fragments a rack: 1.5× the bytes, and a rack plus a node can fail.
  3. 3Durability is a race between disk failures and repair, so repair bandwidth and scrubbing are part of the design.
  4. 4One latest row per key is the commit point: strong read-after-write, prefix listing as a range scan, and safe GC.
C11
    A

    The prompt

    as the interviewer says it
    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.

    B

    Clarifying questions

    ask, then assume
    questionanswer assumedit decides
    How large is an average object?1 MB; many under 64 KBObject count, index size
    Reads and writes a day?2 × 10⁹ PUTs, 10 GETs per PUTFront ends, index shards
    Must a list see a write at once?Yes: strong list-after-writeOne index, range shards
    Overwrite or append?Whole-object overwrite onlyImmutable versions, GC
    What must fail without loss?A rack and one more diskCode, placement
    Keep old versions?Per bucket: on or offGC rule
    One region or many?One region; copy to a second on requestReplication scope

    I ask about object sizes, the read-write mix, listing consistency and the failures to survive. Each answer changes the layout.

    C

    Estimates

    arithmetic shown
    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.

    D

    Eleven nines, worked

    the B9 model, recorded grid
    layoutbytessurvives2%, 24 h5%, 72 h
    3 replicas3×210.28.1
    Reed-Solomon 6+31.5×312.49.4
    Reed-Solomon 10+41.4×415.411.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.

    E

    API

    buckets, keys, versions, parts
    callsuccesserrors
    PUT /{bucket}/{key} + Content-SHA256200 ETag, version-id400 bad digest, 403, 503
    GET /{bucket}/{key}[?version-id] + Range200 or 206 bytes404, 503
    HEAD /{bucket}/{key}200 size, ETag404
    DELETE /{bucket}/{key}204, delete marker403
    GET /{bucket}?prefix=&after=&limit=200 keys + next token403
    POST /{bucket}/{key}?uploads200 upload-id403
    PUT ...?upload-id=&part=n200 part ETag400 bad digest
    POST ...?upload-id= + part list200 ETag, version-id400 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.
    multipart uploadpseudo code
    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
    1. 1Parts live apart from the key. A GET still returns the previous version.
    2. 2A network failure costs one part, not the whole upload. A part with a bad checksum is refused.
    3. 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
    Go: multipart, reused from the B9 labgo
    // 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.

    F

    Data model

    the index, and what a node keeps
    1 key : 1 visible version1 version : k + m fragmentsversionsPKidgrows per writeFKbucket_idkeybyte ordersize_bytes, etagis_deletedelete markercommittedfalse until commitone row per PUT or DELETElatestPKbucket_idPKkeyUQversion_idFK versionscommit = 1 upsert herefragmentsPKversion_idPKidx0 to k+m-1node_idwhich diskchecksumfor scrubk + m rows per versionstorage node(version, idx)file namebytessize / kno index, no key names
    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.
    the metadata index in Postgressql
    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)
    );
    1. 1Byte order, not language order, so a prefix is one contiguous range of the index.
    2. 2A version is written before it is visible. A versioned bucket keeps committed versions; GC takes the rest.
    3. 3Scrub compares each fragment with this. It also catches a disk that returns the wrong block.
    4. 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.

    G

    The architecture

    click a step; its path lights up
    ClientSDK, signedFront endauth, encode, routeWorkersrepair, scrub, GCMetadata indexPostgres, range shardsStorage nodesrack 1 of 6Storage nodesrack 2 of 6Storage nodesrack 6 of 6

    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.

    H

    Capabilities used

    what each part gives you
    toolcapabilitywhat it gives this designalso used for
    Front endReed-Solomon k + mAny k of k + m fragments rebuild the object; 1.5× bytes for 6+3.RAID 6, backups, QR codes
    PlacementAt most ⌈(k+m)/racks⌉ fragments a rackA rack failure takes at most 2 of 9.Replica placement in any database
    Storage nodeA checksum per fragmentA bad fragment is skipped on read and found by scrub.File systems, network frames
    WorkersRepair from any kSpare fragments come back before the next failure.Distributed file systems
    PostgresINSERT ... ON CONFLICT ... WHEREThe commit point: the newest version wins, read-after-write holds.Optimistic writes, idempotency
    PostgresB-tree in byte orderA prefix list is one index range scan.Time ranges, autocomplete
    PostgresForeign keysGC cannot delete a live version.Any referential rule
    PostgresRange shards with sync replicasThe index grows past one machine and survives a primary loss.Large key-value indexes
    CDNEdge cacheHot 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.

    I

    Deep dive: placement and the data path

    where do the 9 fragments go?
    spread: at most 2 per rackrack 1 failsd0n0f6f0d0n3rack 2f5f7d1n2d1n3rack 3d2n0f8d2n2f1rack 4d3n0d3n1f4d3n3rack 5f2d4n1d4n2d4n3rack 6f3d5n1d5n2d5n32 of 9 lost: readablerandom nodes, racks ignoredrack 1 failsf5f7f0f6rack 2d1n0f8d1n2d1n3rack 3d2n0f1d2n2d2n3rack 4d3n0d3n1f2f4rack 5d4n0d4n1f3d4n3rack 6d5n0d5n1d5n2d5n34 of 9 lost: unreadable
    placementstatusrack 1 fails
    6+3 on 9 random nodesNot approved11 of 1,000 lost
    3 replicas on 3 random nodesNot approved3 of 1,000 lost
    6+3, at most 2 a rackApproved0 lost; also with 1 more node
    6+3, 3 a zone over 3 zonesA zone must fail safelyEvery 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 and getpseudo code
    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)
    1. 1The 6 data fragments are the object cut in 6. A read with all of them up needs no decoding.
    2. 2The rule prefers the rack with the fewest fragments, and skips down nodes.
    3. 3Fragments first, commit last. A crash in between leaves garbage, never a broken object.
    4. 4Reads of the index go to the primary of its shard. A replica could return the version before the commit.
    5. 5A degraded read. In the lab, one node down made 267 of 1,000 reads degraded and none failed.
    6. 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
    Go: placego
    // 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
    }
    
    Go: put and commitgo
    // 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
    }
    
    Go: getgo
    // 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
    }
    
    Go: Reed-Solomon rebuild, reused from the B9 labgo
    // 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.

    J

    Deep dive: repair and scrubbing

    durability is a race
    methodstatuscost
    Rebuild only when an object is readNot approvedCold objects lose spares in silence.
    Rebuild every object on a dead node, from any 6Approved3.3 B read per B rebuilt
    3 replicas: copy a survivorHot or small data0.9 B read per B rebuilt
    Wait until 2 of 3 spares are goneRepair bandwidth is shortLess traffic; one spare left at times
    Scrub every disk every few weeksApprovedReads 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 and scrubpseudo code
    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
    1. 1This must stay at zero. A lost object is a page to the on-call engineer, not a log line.
    2. 2One read of 6 fragments rebuilds all the missing ones of that object.
    3. 3The same rule as a write, so the repaired object again survives a rack.
    4. 4A disk can return bad bytes with no error. Scrub finds them while spares still exist.
    5. 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
    Go: repairgo
    // 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
    }
    
    Go: scrubgo
    // 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
    }
    
    Go: GC on the nodesgo
    // 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
    }
    
    Go: durability model, reused from the B9 labgo
    // 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.

    K

    Deep dive: the metadata index

    consistency, listing, garbage
    indexstatusprefix list
    One PostgresUp to a few TB of rowsOne range scan
    Hash-sharded key-value storeNot approvedAsks every shard
    Range shards by (bucket, key)ApprovedOne shard or neighbours
    Consensus-replicated key rangesMany regionsOne range, quorum writes
    Eventually consistent replicasNot approvedCan 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, list, gcpseudo code
    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
    1. 1Two writers commit in any order; the larger version id stays. The other version is garbage.
    2. 2prefix_end is the prefix with its last byte raised by one. The scan touches only matching keys.
    3. 3A key token, not an offset, so new keys do not shift the pages.
    4. 4Longer than the slowest write, so GC never takes a version between write and commit.
    5. 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
    SQL: commitsql
    -- 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;
    SQL: listsql
    -- 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;
    SQL: garbagesql
    -- 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;
    Go: begin and commitgo
    // 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
    }
    
    Go: listgo
    // 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
    }
    
    Go: garbagego
    // 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.

    L

    Try it: fail racks, then repair

    recorded from the lab cluster

    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.

    rack 1
    d0n0392
    d0n1f6389
    d0n2f0362
    d0n3331
    rack 2
    d1n0f5364
    d1n1f7401
    d1n2389
    d1n3373
    rack 3
    d2n0399
    d2n1f8373
    d2n2362
    d2n3f1374
    rack 4
    d3n0380
    d3n1350
    d3n2f4369
    d3n3359
    rack 5
    d4n0f2370
    d4n1379
    d4n2386
    d4n3387
    rack 6
    d5n0f3383
    d5n1402
    d5n2363
    d5n3363

    bars: fragments each node holdsf0 to f8: object photos/2026/0172.jpg, 0 of 9 fragments gone, readable

    objects by fragments gone
    01000
    10
    20
    30
    40
    50
    60
    70
    80
    90
    1,000 of 1,000 objects readable, 0 lost. 0 reads needed a parity fragment.

    Spread placement keeps every object readable through a rack and a node. Random placement lost 47 of 1,000 objects to the same failures.

    M

    Failure cases

    what breaks, and why it stays correct
    eventresultwhy it is safesaved by
    A disk or node diesReads of its fragments become degraded.Any 6 of 9 fragments rebuild the object. Repair restores the spares.Erasure code
    A rack loses powerUp to 2 fragments of each object are gone.7 remain, one more than needed.Placement
    A disk returns bad bytesOne fragment fails its checksum.The read skips it; scrub deletes it; repair rebuilds it.Checksums
    The front end dies before commit9 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 onceTwo versions.The commit keeps the larger version id; the other becomes garbage.Conditional upsert
    An index primary failsWrites to its key range pause.A synchronous replica has every commit and takes over in seconds.Sync replica
    A batch of disks fails togetherFailures are no longer independent.Mix disk models and ages across racks; the model assumes independence.Operations
    One object gets most of the readsIts 6 data nodes run hot.Serve it from a CDN, or read different sets of 6 of its 9 fragments.CDN
    N

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1One server23 replicas36+3, racks4+ index shards5+ cellsmore load →
    Index capacity against demand1k10k100k1M10M100MDemand, PUT peak: 69,444 operations per secondDemand, PUT peak69,444Demand, GET peak: 694,444 operations per secondDemand, GET peak694,444One shard, commits: 5,000 operations per secondOne shard, commits5,000One shard, HEADs: 80,000 operations per secondOne shard, HEADs80,000250 shards, commits: 1,250,000 operations per second250 shards, commits1,250,000250 shards, HEADs: 20,000,000 operations per second250 shards, HEADs20,000,000operations per second, log scale
    stepaddit handlesmove up when you see
    1One 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.
    23 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.
    3Erasure 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.
    4Range-sharded index, split by (bucket, key).About 250 shards for 10¹² objects.One cluster's repair traffic, network or racks hit a limit.
    5Cells: 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.

    O

    Follow-ups

    what the interviewer asks next
    questionanswergo 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
    P

    Drill

    predict, then reveal

    0 of 9 known

    1. Why Reed-Solomon 6+3 and not 3 replicas?

    2. There are 9 fragments and 6 racks. Why at most 2 in a rack?

    3. A front end wrote all 9 fragments, then crashed before the commit. What do readers see?

    4. Why does repair of a 6+3 object move more bytes than repair of a replica?

    5. Eleven nines, 10¹² objects. How many objects are lost a year?

    6. Why does repair speed decide durability?

    7. Why shard the index by key range and not by hash?

    8. How does GC avoid deleting an object a reader can see?

    9. What does scrub find that repair does not?

    Q

    Numbers to say

    measured, derived or cited
    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.