Skip to main content

Dropbox (Cloud Storage & Sync)

Sharpened prompt. Design a file sync service for 700M users and 500 PB of data, where a one-byte edit to a 4 GB video uploads in under a second, the same file stored by 10,000 users occupies space once, two devices editing offline never silently lose data, and renaming a folder containing a million files is an O(1) operation.

The naive design — upload the file, download the file — is wrong in a specific and interesting way: the unit of work is not the file, it is the chunk. Every good property in this system falls out of that one decision.

1. Problem framing​

Functional requirements​

  • Upload, download, and delete files; sync across all of a user's devices.
  • Efficient sync: transfer only what changed.
  • Shared folders with permissions; version history and restore.
  • Offline editing with automatic conflict handling.

Non-functional requirements​

PropertyTargetConsequence
Bandwidth efficiencySmall edits transfer in KB, not GBContent-defined chunking, not whole-file upload
Storage efficiencyGlobal dedup across all usersContent-addressable storage keyed by hash
Sync latencyChange visible on other devices in < 3sPush notifications, never polling
Durability11 ninesErasure-coded object storage
ConsistencyEventual, with no silent data lossConflicts are surfaced, never resolved by deletion

Back-of-the-envelope​

Users: 700M, ~100M daily active, ~3 devices each = 300M sync clients
Data: 500 PB logical; dedup + compression ~ 40% saving -> ~300 PB physical
Files: 700M users × 10,000 files = 7 trillion file-version rows (metadata is the
harder database problem, not the bytes)
Chunks: 4 MB average -> 500 PB / 4 MB = 125 billion chunks
Chunk index: 125B × (32 B hash + 24 B location) = 7 TB of index -> sharded KV store
Notifications: 300M long-lived connections; at 10 KB/conn of kernel + app state = 3 TB RAM
across the notification tier -> ~1,500 nodes at 200k connections each

Two independent scaling problems live here: bytes (petabytes, solved by object storage and dedup) and metadata (trillions of small mutable rows, solved by sharding). Separating them explicitly is the first good move.

2. High-level architecture​

The six-step numbering is the whole protocol and it is worth drawing exactly like this: ask what's missing, upload only that, commit the manifest atomically, notify, let others pull. Commit-after-upload is what makes the system crash-safe — a manifest never references a chunk that is not durably stored.

3. Component inventory​

ComponentConcrete choiceWhy this one
ChunkerFastCDC (Gear-hash rolling)3–10× faster than Rabin fingerprinting for identical dedup quality
Chunk hashBLAKE3 or SHA-256BLAKE3 is ~10× faster and parallelises; SHA-256 if you need FIPS
Chunk indexScyllaDB / Cassandra125B rows, point lookups by hash, high write rate, no joins ever
Object storageS3 (or an in-house Magic Pocket) with Reed–Solomon erasure codingRS(10,4) gives 1.4× overhead vs 3× for replication — at 300 PB that is the whole storage budget
Metadata DBSharded MySQL/Postgres by namespace_idNeeds transactions per namespace; a namespace never spans shards
NotificationsWebSocket tier + Redis Pub/Sub300M idle connections; Go or Erlang for cheap goroutines/processes
Client storeSQLiteEvery client is a small database; crash-consistent local state is mandatory

4. Data model​

-- Sharded by namespace_id. A namespace = a user's root or a shared folder.
CREATE TABLE files (
namespace_id BIGINT,
path VARBINARY(4096),
file_id BIGINT,
latest_rev BIGINT,
PRIMARY KEY (namespace_id, path)
);

CREATE TABLE revisions (
file_id BIGINT,
rev BIGINT,
manifest_id BINARY(32), -- hash of the ordered chunk list
size BIGINT,
modified_by BIGINT, -- device id
vector_clock BLOB, -- for concurrent-edit detection
created_at TIMESTAMP,
PRIMARY KEY (file_id, rev DESC)
);

-- Global, not sharded by user: this is where dedup happens.
CREATE TABLE chunks (
chunk_hash BINARY(32) PRIMARY KEY,
size INT,
location TEXT, -- object key + offset
refcount COUNTER
);

The manifest — an ordered list of chunk hashes — is the pivot of the whole design. A file is its manifest; the bytes are shared infrastructure.

5. The toughest parts​

5.1 Delta sync: why fixed-size chunking fails​

Why it's hard. Split a file into fixed 4 MB blocks and insert one byte at the front. Every subsequent block shifts by one byte, so every block hash changes and you re-upload the entire 4 GB file for a one-byte edit. This is the "boundary shift" problem, and it makes fixed chunking useless for anything but append-only files.

Original: [--- A ---][--- B ---][--- C ---]
Insert 1 byte at offset 0:
Fixed: [x-- A' --][-- B' ---][-- C' --] all 3 hashes differ -> re-upload all
CDC: [x][--- A ---][--- B ---][--- C ---] only the first boundary moved

Solution — content-defined chunking. Slide a window over the data, compute a rolling hash, and cut a boundary wherever the hash matches a mask. Because boundaries depend on content, an insertion shifts only the chunks it physically touches; everything downstream re-aligns.

# FastCDC: Gear hash + a mask. Cheaper than Rabin, same dedup ratio.
GEAR = [random.getrandbits(64) for _ in range(256)]
MASK_S, MASK_L = 0x0003590703530000, 0x0000d90003530000 # normalised chunking

def chunk(data: bytes, min_sz=2<<20, avg_sz=4<<20, max_sz=8<<20):
i, h, start = min_sz, 0, 0
while i < len(data):
h = ((h << 1) + GEAR[data[i]]) & 0xFFFFFFFFFFFFFFFF
# Two masks: strict below the average (resist small chunks),
# loose above it (force a cut before max). This tightens the size distribution.
mask = MASK_S if (i - start) < avg_sz else MASK_L
if (h & mask) == 0 or (i - start) >= max_sz:
yield data[start:i]
start, i, h = i, i + min_sz, 0
continue
i += 1
if start < len(data):
yield data[start:]

The min/avg/max triple matters: without a minimum you get pathological 12-byte chunks that explode the index, and without a maximum a low-entropy region (a run of zeros) never triggers a boundary.

Then the upload protocol becomes trivially efficient:

chunks = list(chunk(open(path,'rb').read()))
hashes = [blake3(c).digest() for c in chunks]
missing = await api.probe(hashes) # one round trip: which do you lack?
await asyncio.gather(*(api.put_chunk(h, c) # upload only the diff
for h, c in zip(hashes, chunks) if h in missing))
await api.commit(path, manifest=hashes) # atomic: the file now exists at this rev

For a one-byte edit in a 4 GB file, missing contains one or two chunks. 8 MB transferred instead of 4 GB.

5.2 Global deduplication and the security problem hiding inside it​

Why it's hard. Dedup is enormously valuable — one copy of a popular installer instead of 10,000 — but it creates a cross-user side channel. If the client asks "do you already have chunk H?" and skips the upload when the answer is yes, an attacker who can guess a file's contents can confirm its existence on the service. With a low-entropy file (a form with a known template and an unknown salary field) that is a practical extraction attack.

Solution — content-addressable storage with scoped dedup and reference counting.

Chunk store (global, immutable, keyed by hash):
blake3(chunk) -> { object_key, size, refcount }

Manifest (per file version):
[h1, h2, h3, ...] ordered list

Dedup scope:
- Within one user / one team: full client-side dedup ("do you have H?")
- Across users: server-side only. The client always uploads;
the server discards the duplicate after receiving it.

The asymmetry is the fix: you keep the storage saving (the server still stores one copy) while removing the side channel (the client's bandwidth no longer reveals what the server knows). State this trade explicitly — "cross-user client-side dedup saves bandwidth but leaks existence, so I keep dedup server-side across trust boundaries and client-side within them."

Reference counting is the other half, and it is genuinely error-prone. A naive decrement-then-delete races with a concurrent upload that just incremented from zero. The standard fix is delayed garbage collection: never delete on refcount 0; instead mark the chunk with a tombstone timestamp and have a GC job delete only chunks whose refcount has been 0 for, say, 30 days, re-verifying atomically at delete time. Chunks are cheap; a chunk deleted while someone references it is unrecoverable data loss.

Also worth naming: convergent encryption (encrypt each chunk with a key derived from its own hash) lets you dedup encrypted data, at the cost of the same confirmation-of-file attack. Most providers now dedup only within a trust boundary for exactly this reason.

5.3 Two devices, one file, no network: conflict resolution​

Why it's hard. A laptop and a phone both edit notes.txt while offline. Both come back online. Last-writer-wins by wall clock is wrong — clocks are skewed, and the "loser" silently loses work. For arbitrary binary files (a Photoshop document, a video) there is no meaningful merge. And you cannot ask the user about every conflict; sync must be automatic.

Solution — detect concurrency with vector clocks, then preserve both branches.

class VectorClock(dict): # device_id -> counter
def happens_before(self, other) -> bool:
return (all(self.get(d, 0) <= other.get(d, 0) for d in set(self) | set(other))
and self != other)

def concurrent_with(self, other) -> bool:
return not self.happens_before(other) and not other.happens_before(self)

def merge(local: Rev, remote: Rev) -> Action:
if local.vc.happens_before(remote.vc): return FastForward(remote) # simple update
if remote.vc.happens_before(local.vc): return KeepLocal() # we're ahead
# Genuinely concurrent: never pick a winner for opaque bytes.
return CreateConflictCopy(
keep=remote, # canonical path takes the remote
rename_local=f"{local.stem} (conflicted copy from {local.device_name}) "
f"{local.mtime:%Y-%m-%d}{local.suffix}")

Vector clocks detect concurrency correctly where timestamps only guess. The clock is per-file and pruned to the devices that actually touched it, so it stays small (a few entries) rather than growing with the account's device count.

Two refinements worth mentioning:

  • Type-aware merging. For formats you understand — plain text, and increasingly structured formats — a three-way merge against the common ancestor (which the manifest history gives you for free) resolves the majority of real conflicts silently. Fall back to a conflicted copy only when the merge fails.
  • Directory conflicts are harder than file conflicts. Device A moves /work/report.doc to /archive/; device B edits it in place. Neither operation conflicts at the file level, but applying them in different orders on different devices yields divergent trees. Model moves as (delete, create) on the namespace with their own vector clock entry, and resolve namespace operations before content operations.

5.4 Telling 300 million devices that something changed​

Why it's hard. Polling is the obvious design and it is catastrophic at this scale: 300M clients polling every 30 seconds is 10M QPS of "anything new?" requests, ~99.9% of which return "no." On mobile it also destroys battery, because each poll wakes the radio.

Solution — split the notification path from the data path, and keep the notification tiny.

Client Notification tier Metadata tier
│ WS connect (account_id, cursor) │ │
├──────────────────────────────────►│ │
│ │ SUBSCRIBE acct:{shard}
│ ├──────────────────►│ (Redis Pub/Sub)
│ │ │
│ [ idle, ~10 KB of state ] │ commit on device B
│ │◄──────────────────┤
│ {"changed": true, "cursor": 91} │ │
│◄──────────────────────────────────┤ │
├── GET /delta?cursor=88 ───────────────────────────────►│ (normal HTTP, load balanced)

The push carries no data — just "you are stale, here is the new cursor." The client then pulls the delta over ordinary, cacheable, load-balanced HTTP. This keeps the stateful connection tier dumb and cheap, and means a notification storm cannot corrupt anything; worst case, clients pull.

Implementation notes that matter: use Go or Erlang for the connection tier (300M connections at ~10 KB each needs cheap concurrency, not threads); shard Pub/Sub channels by account so no node subscribes to everything; fall back to long polling where WebSockets are blocked, and to FCM/APNs silent push when the app is backgrounded on mobile — the OS will not let you hold a socket.

The cursor is a per-namespace monotonic sequence number. GET /delta?cursor=N returns everything after N, which makes the client's catch-up logic identical whether it was offline for 2 seconds or 2 months.

5.5 Renaming a folder with a million files​

Why it's hard. If file paths are stored as strings, mv /Projects /Archive/Projects requires updating a million rows — a transaction that locks a shard for minutes, replicates as a million binlog events, and generates a million sync notifications, each of which makes clients re-download nothing but re-write their local index.

Solution — store the tree structurally, not as flattened paths. Each file row references a parent_id; the full path is derived. A rename becomes a single-row update.

-- One row changes. Every descendant's path is derived, so nothing else moves.
UPDATE nodes SET parent_id = :new_parent, name = :new_name WHERE node_id = :folder_id;

The cost lands on reads: resolving /a/b/c/d.txt needs four lookups. Fix that with a materialised path cache on each node (an ancestry array like [root, 12, 45, 91], which supports a prefix query for "everything under 45") refreshed lazily, or a closure table if you need strict SQL ancestry queries.

The sync side needs care too. Clients must be told "subtree 45 moved" as one event, not a million file events. That means the delta stream carries namespace operations as first-class entries, and the client applies a local subtree move — a directory rename on the local filesystem, touching no file contents.

Shared folders make this sharper. A shared folder is its own namespace_id and lives on its own shard, mounted into each member's tree by reference. Moving a file into a shared folder is therefore a cross-namespace operation: it changes who can see the bytes, so it is a copy-and-delete with a permission re-evaluation, not a pointer update. Recognising that "move between namespaces is not a move" is a strong detail.

5.6 Resumable uploads over a hostile network​

Why it's hard. A 4 GB upload from a phone on a train will be interrupted. Restarting from zero is unacceptable. The upload also has to survive the server restarting mid-transfer, and it must not leave half-written files visible to other devices.

Solution — chunk-level idempotent uploads plus a two-phase commit on the manifest.

POST /2/files/upload_session/start
→ { "session_id": "s-8f14e" }

POST /2/files/upload_session/append?session_id=s-8f14e&offset=12582912
Content-Length: 4194304
Idempotency-Key: s-8f14e:12582912 ; retrying this exact call is a no-op
→ 200 { "offset": 16777216 }

POST /2/files/upload_session/finish
{ "session_id": "s-8f14e", "manifest": ["h1","h2",...], "path": "/video.mp4" }
→ 200 { "rev": 4471 } ; the file becomes visible only here

Three properties fall out. Each chunk PUT is idempotent and independent, so they can run in parallel and retry freely (chunks are content-addressed — writing the same hash twice is definitionally a no-op). The file is invisible until finish, so no other device ever sees a partial file. And on resume, the client asks the server which chunks it already has and skips them — the same probe call from 5.1, reused.

Add adaptive parallelism (4–16 concurrent chunk uploads, tuned by observed throughput) and per-chunk checksums verified server-side, because silent corruption over a bad connection is real and content-addressing makes it free to detect: if the received bytes do not hash to the claimed name, reject.

5.7 Storage economics at 300 petabytes​

Why it's hard. At this scale, storage strategy is the business model. 3× replication of 300 PB is 900 PB of disk. Most data is cold — the overwhelming majority of files are never read after the first week — but "cold" is a prediction, and getting it wrong means a user waits 12 hours for their file.

Solution — erasure coding plus temperature tiering, with hot metadata always fast.

  • Reed–Solomon(10,4): split each chunk into 10 data and 4 parity fragments across failure domains. Any 10 of 14 reconstruct the chunk. Overhead is 1.4× instead of 3×, tolerating 4 simultaneous failures instead of 2. At 300 PB that is roughly 480 PB of disk instead of 900 PB — the single largest cost decision in the design.
  • Tiering by access recency: hot on SSD/standard object storage, warm on infrequent-access, cold on Glacier-class. Move on a rolling 30/90-day no-access rule, and never tier metadata — the file must always appear instantly in the UI even if the bytes take time.
  • Cross-region replication for the metadata, single-region-plus-EC for the bytes, unless the customer pays for geo-redundancy. Be explicit that this is a durability/cost trade the product makes, not an accident.

The reconstruction cost is the trade to name: reading a chunk whose fragment is missing requires fetching 10 fragments and doing the RS math, which adds latency and network load. That is why the hot tier stays replicated and only warm and cold data is erasure-coded.

6. What breaks first​

EventFirst failureMitigation
A user syncs a 500k-file folderMetadata shard write saturationBatch commits; rate-limit per namespace; chunked delta application
Viral shared folder (1M members)Pub/Sub fan-out on one channelShard the notification channel; hierarchical fan-out for large namespaces
Chunk index hot partitionScyllaDB shard saturation on a popular chunkHash-based partitioning is uniform by construction; cache popular chunk locations
Client clock badly wrongBogus conflict detectionVector clocks, never wall-clock comparison
Mass file deletionGC deletes chunks still referenced elsewhereDelayed GC with re-verification; tombstones, never immediate delete
Notification tier node loss200k clients reconnect at onceRandomised reconnect backoff; clients function offline and resync via cursor
Region outageMetadata unavailableMetadata replicated cross-region; bytes served from surviving EC fragments

7. Cheat sheet​

  • Chunk, don't file. FastCDC with min/avg/max sizes; BLAKE3 hashes; a file is an ordered manifest of chunk hashes.
  • Protocol: probe → upload missing → atomic manifest commit → notify → peers pull delta. Commit last.
  • Dedup: content-addressable storage with refcounts; client-side dedup only inside a trust boundary; delayed GC.
  • Conflicts: vector clocks detect concurrency; concurrent binary edits become conflicted copies. Never last-writer-wins.
  • Notify: WebSocket push carrying only a cursor; data pulled over plain HTTP. Never poll.
  • Namespace: parent_id tree so a folder rename is one row; deltas carry subtree moves as single events.
  • Storage: RS(10,4) at 1.4× instead of 3× replication; tier by access age; never tier metadata.
  • The one-liner: "Two systems in a trench coat — an immutable content-addressed blob store that dedups globally, and a mutable per-namespace metadata tree that everything else is really about."