Problem and scope
A file sync service keeps a folder of files identical across all of a user's devices, and stores a copy in the cloud. You save a document on your laptop, and a few seconds later it appears on your phone. You share a folder with a teammate, and their edits show up on your machine. Dropbox, Google Drive, Microsoft OneDrive and iCloud Drive all solve this problem.
It looks like "upload and download files", but the interesting parts are hidden:
- How do you avoid re-uploading a 2 GB file when the user changed one paragraph?
- How do other devices find out, within seconds, that something changed?
- What happens when two people edit the same file while one of them is offline?
- How do you store hundreds of petabytes without storing the same bytes many times?
Interviewers use this problem to test whether you can separate metadata (names, folders, versions, permissions) from content (the bytes), and whether you can reason about synchronisation and conflicts rather than just storage.
Clarifying questions
| Question | Assumed answer |
|---|---|
| Desktop sync client, mobile app, web, or all? | All three; the desktop client does continuous sync |
| Maximum file size? | 50 GB |
| Real-time collaborative editing (like Google Docs)? | No. We sync whole files; collaborative editing is a different design |
| Sharing with other users? | Yes: folders and files, view or edit permission |
| Version history? | Yes, keep previous versions for 30 days |
| Offline edits? | Yes, the desktop client works offline and syncs later |
| How fast should changes reach other devices? | A few seconds when online |
| Storage quota per user? | Yes, for example 15 GB free, more on paid plans |
| End-to-end encryption? | Not required; mention the trade-off |
The answer to "real-time collaborative editing" is the most important one. Character-by-character co-editing needs operational transforms or CRDTs (data structures that merge concurrent edits automatically). File sync deals in whole files and versions, which is much simpler and what we design here.
Functional and non-functional requirements
Functional:
- Upload, download, rename, move and delete files and folders.
- Automatically sync changes to all of the user's devices.
- Share a file or folder with others with view or edit rights.
- Keep a version history and allow restoring old versions.
- Support offline edits and resolve conflicts when devices reconnect.
- Resume interrupted uploads and downloads.
Non-functional:
- Durability above everything. Losing a user's file is the worst possible failure. Target the durability of a good object store (it is commonly quoted as "eleven nines" annual object durability for major cloud stores).
- Consistency of metadata. Every device must eventually see the same folder tree, and no edit may be silently lost.
- Low sync latency: a few seconds from save to appearing elsewhere.
- Bandwidth efficiency: only send what changed. Users on mobile data and slow broadband notice waste.
- Scalability: hundreds of millions of users, hundreds of petabytes.
- Availability: the service should be readable even if some writes are delayed.
Back-of-the-envelope estimates
Assumptions:
- 100 million registered users, 20 million daily active users.
- Average stored data per user: 2 GB.
- Each daily active user changes 10 files per day, average file size 1 MB.
- Average 2,000 files per user.
- About 1.5 connected devices per daily active user.
Storage. 100 million users × 2 GB = 200 PB of logical data. If deduplication saves about 25% (an assumption; the real figure depends heavily on the user base), we store 150 PB of unique bytes. Durable storage with erasure coding at roughly 1.5× overhead brings the raw disk to about 225 PB. Version history adds more on top.
Write rate. 20 million users × 10 changes = 200 million file changes per day. That is 200,000,000 ÷ 86,400 ≈ 2,300 changes per second on average. Working hours cluster activity, so plan for about three times that at peak: roughly 7,000 per second.
Upload bandwidth. In the worst case each change uploads the whole 1 MB file: 200 million × 1 MB = 200 TB per day, about 18.5 Gbps on average. Chunk-level deltas (explained below) reduce this a lot for large files that change slightly.
Metadata. 100 million users × 2,000 files = 200 billion file entries. At around 500 bytes per entry (path, ids, sizes, hashes, timestamps), that is 100 TB of metadata. Much smaller than the content, but far too large for one database, and it is the part that needs strong consistency. It must be sharded.
Connections. 20 million daily users × 1.5 devices ≈ 30 million devices, many of which stay connected for change notifications. Even if only a third are connected at the same moment, that is about 10 million open long-lived connections.
Chunks. With 4 MB chunks, a 2 GB user has about 500 chunks. A 10 MiB file splits into 4 MiB + 4 MiB + 2 MiB, three chunks.
Interview tip
Point out the split early: "Content is 200 PB but simple, immutable blobs. Metadata is only about 100 TB but it is where consistency, conflicts and permissions live. I will store them in different systems and spend my design time on the metadata and sync protocol."
API design
The client talks to two services. The metadata service handles names, folders, versions and permissions. The block service handles raw chunk bytes.
# Block service: content-addressed chunks
POST /v1/blocks/check
body: { "hashes": ["9f2c...", "a71b...", "03de..."] }
200 -> { "missing": ["a71b..."] }
PUT /v1/blocks/a71b... (raw chunk bytes, <= 4 MB)
201 -> {}
GET /v1/blocks/a71b...
200 -> raw bytes
# Metadata service: files, versions, sync
POST /v1/files/commit
body: { "namespaceId": "ns_42", "path": "/Reports/q3.xlsx",
"blockHashes": ["9f2c...", "a71b...", "03de..."],
"size": 9437184, "parentRev": 17,
"clientMtime": "2026-10-10T09:12:00Z" }
200 -> { "fileId": "f_88", "rev": 18 }
409 -> { "error": "CONFLICT", "currentRev": 18 }
GET /v1/changes?namespaceId=ns_42&cursor=c_1043
200 -> { "entries": [ ... ], "cursor": "c_1051", "hasMore": false }
GET /v1/changes/longpoll?cursor=c_1051&timeout=60
200 -> { "changes": true } or after 60 s { "changes": false }
GET /v1/files/f_88/revisions
POST /v1/files/f_88/restore body: { "rev": 16 }
POST /v1/shares
body: { "path": "/Reports", "grantee": "user_9", "role": "EDITOR" }
Key ideas in this API:
- Blocks are addressed by their hash. A block's id is the SHA-256 of its bytes. Two identical chunks anywhere in the system have the same id, which gives deduplication for free.
- Check before upload. The client asks which hashes the server is missing and uploads only those.
- Commit is separate from upload. Uploading blocks does not change any file. A file changes only when the client commits a new list of block hashes. This makes the file update atomic: other devices see either the old version or the new one, never half.
parentRevgives optimistic concurrency. The client says "I edited revision 17". If the server's latest is already 18, the commit is rejected with409 Conflict, and the client handles the conflict.- A cursor-based change feed lets each device ask "what changed since the last cursor I saw?"
Data model and storage choice
Content: block storage. Chunks go into object storage (or an in-house blob store), keyed by hash. They are immutable: a chunk with a given hash never changes, so it can be cached and replicated freely.
Metadata: a sharded, strongly consistent database. A sharded relational database (for example MySQL or PostgreSQL shards) or a distributed SQL database. We need transactions so that a commit updates the file row, writes a revision and appends to the change log atomically.
-- A namespace is a sync root: a user's home folder or a shared folder.
CREATE TABLE namespaces (
namespace_id BIGINT PRIMARY KEY,
owner_id BIGINT NOT NULL,
kind VARCHAR(10) NOT NULL -- HOME, SHARED
);
CREATE TABLE files (
namespace_id BIGINT NOT NULL,
file_id BIGINT NOT NULL,
parent_id BIGINT, -- folder containing it
name VARCHAR(255) NOT NULL,
is_folder BOOLEAN NOT NULL,
latest_rev BIGINT NOT NULL,
is_deleted BOOLEAN NOT NULL DEFAULT FALSE,
PRIMARY KEY (namespace_id, file_id),
UNIQUE (namespace_id, parent_id, name)
);
CREATE TABLE revisions (
namespace_id BIGINT NOT NULL,
file_id BIGINT NOT NULL,
rev BIGINT NOT NULL,
size_bytes BIGINT NOT NULL,
block_hashes TEXT NOT NULL, -- ordered list of SHA-256 hashes
author_device BIGINT NOT NULL,
created_at TIMESTAMP NOT NULL,
PRIMARY KEY (namespace_id, file_id, rev)
);
-- Append-only journal; the cursor is the journal position.
CREATE TABLE change_journal (
namespace_id BIGINT NOT NULL,
seq BIGINT NOT NULL, -- increases by 1 per change
file_id BIGINT NOT NULL,
rev BIGINT NOT NULL,
op VARCHAR(8) NOT NULL, -- ADD, MODIFY, MOVE, DELETE
PRIMARY KEY (namespace_id, seq)
);
CREATE TABLE acl (
namespace_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
role VARCHAR(8) NOT NULL, -- OWNER, EDITOR, VIEWER
PRIMARY KEY (namespace_id, user_id)
);
CREATE TABLE block_refs (
block_hash CHAR(64) PRIMARY KEY,
ref_count BIGINT NOT NULL,
size_bytes INT NOT NULL
);
Why shard by namespace? Almost every operation is scoped to one namespace: list a folder, commit a file, read the change journal. Putting a whole namespace on one shard keeps these operations single-shard transactions. A shared folder is its own namespace, so all its members read the same journal on the same shard.
Why a journal per namespace? The seq column gives a strict order of changes inside a namespace. A device's cursor is simply "the highest seq I have applied", so "what changed since cursor 1043" is a cheap range scan on the primary key.
Paths versus ids
Store a parent id and a name, not full path strings. Renaming a folder with 10,000 files inside then changes one row, not 10,000. The full path is computed by walking parents, and cached on the client.
High-level design
+-------------------+ +-------------------+
| Desktop client | | Mobile / web |
| - file watcher | | client |
| - chunker+hasher | +---------+---------+
| - local index DB | |
+----+---------+----+ |
| | |
| blocks | metadata, changes |
v v v
+----------+ +---------------------------------+
| Block | | API gateway (auth, rate limits) |
| service | +----------------+----------------+
+----+-----+ |
| +--------+---------+
v v v
+----------+ +-------------+ +---------------+
| Object | | Metadata | | Notification |
| storage | | service |--->| service |
| (chunks) | +------+------+ | (long poll / |
+----------+ | | WebSocket) |
v +-------+-------+
+------------------+ ^
| Sharded metadata | |
| DB (by namespace)|---------+
| files, revisions,| journal events
| journal, ACLs |
+------------------+
|
v
+------------------+
| Async workers: |
| GC of blocks, |
| previews, search |
+------------------+
The client is a big part of the design. A desktop client has:
- A file watcher that receives operating-system events when files change (inotify on Linux, FSEvents on macOS, ReadDirectoryChangesW on Windows).
- A chunker and hasher that splits changed files and computes hashes.
- A local index, a small database (often SQLite) recording, for each synced file, the revision and block hashes the client last synced. Comparing the disk to this index tells the client what changed while it was offline.
Request flows step by step
Flow 1: a user edits a file on their laptop
- The file watcher reports that
/Reports/q3.xlsxchanged. The client waits briefly (debounces) in case the application is still writing. - The client splits the file into chunks and hashes each one: say hashes
[9f2c, a71b, 03de]. The local index says the previous version was[9f2c, 55e0, 03de], so only the middle chunk is new. - The client calls
POST /v1/blocks/checkwith the three hashes. The server replies that onlya71bis missing. - The client uploads chunk
a71bwithPUT /v1/blocks/a71b. The block service verifies the bytes hash toa71band stores them. - The client calls
POST /v1/files/commitwith the new hash list andparentRev: 17. - The metadata service, in one transaction on the namespace's shard: checks that the latest revision is still 17, inserts revision 18, updates
files.latest_rev, increments block reference counts, and appendsseq 1051to the change journal. - After commit, it publishes "namespace ns_42 has new changes" to the notification service.
- The client updates its local index to revision 18.
Flow 2: the user's phone picks up the change
- The phone holds a long-poll or WebSocket connection to the notification service, subscribed to namespace
ns_42with cursorc_1043. - The notification service signals "changes available" (it does not send file contents).
- The phone calls
GET /v1/changes?cursor=c_1043and receives the entries up toseq 1051, including filef_88at revision 18 with its block list. - The phone compares the block list with chunks it already has locally. A phone typically syncs lazily: it may only update the listing, and download the chunks when the user opens the file.
- To download, it fetches only the missing chunks via
GET /v1/blocks/{hash}, reassembles the file, verifies the hashes, and stores the new cursorc_1051.
Deep dive 1: chunking and deduplication
Why chunk at all?
Chunking means splitting a file into pieces and treating each piece as a separate stored object. It gives four benefits:
- Delta sync: only changed chunks are uploaded.
- Deduplication: identical chunks across files and users are stored once.
- Resumable transfers: a failed transfer restarts at the last chunk, not at byte zero.
- Parallelism: many chunks can upload at once.
Fixed-size chunking
Cut the file every 4 MB. It is simple and fast. The problem is the boundary shift: insert one byte at the start of the file and every following byte moves by one position. Every chunk now contains different bytes, every hash changes, and the client re-uploads the whole file even though the change was tiny.
original: [AAAA][BBBB][CCCC][DDDD]
insert X: [XAAA][ABBB][BCCC][CDDD][D]
every chunk changed -> upload everything
Fixed-size chunking still works well for edits that overwrite bytes in place, such as databases and virtual machine images, and for appends at the end.
Content-defined chunking (CDC)
Content-defined chunking chooses chunk boundaries based on the content itself, not on positions. A rolling hash is computed over a small sliding window of bytes (for example 48 bytes). Whenever the hash matches a pattern, such as "the lowest 12 bits are zero", the chunker cuts there. Because the cut points depend only on nearby bytes, an insertion only changes the chunk where it happened; later boundaries move along with the content and line up again.
cuts (*) fall where the content matches a pattern
original: [AAAAA*][BBBBBBB*][CCC*][DDDDD]
insert X: [XAAAAA*][BBBBBBB*][CCC*][DDDDD]
only the first chunk changed
With "low 12 bits zero", a cut happens on average once every 2¹² = 4,096 bytes. Real systems add a minimum and maximum chunk size so that unlucky content does not produce tiny or huge chunks. Well-known rolling hashes include Rabin fingerprints and the faster Gear hash used by FastCDC.
The following Python program demonstrates the difference. It builds a 200,000-byte file, inserts 5 bytes at the front, and counts how many chunks would have to be uploaded again.
import hashlib
import random
WINDOW = 48 # bytes in the rolling window
MASK = (1 << 12) - 1 # boundary when low 12 bits are zero: ~4 KiB average
MIN_SIZE, MAX_SIZE = 1024, 16384
BASE, MOD = 257, (1 << 61) - 1
POW = pow(BASE, WINDOW, MOD)
def cdc_chunks(data):
chunks, start, h = [], 0, 0
for i, b in enumerate(data):
h = (h * BASE + b) % MOD
if i - start >= WINDOW:
h = (h - data[i - WINDOW] * POW) % MOD
size = i - start + 1
if (size >= MIN_SIZE and (h & MASK) == 0) or size >= MAX_SIZE:
chunks.append(data[start:i + 1])
start, h = i + 1, 0
if start < len(data):
chunks.append(data[start:])
return chunks
def fixed_chunks(data, size=4096):
return [data[i:i + size] for i in range(0, len(data), size)]
def digests(chunks):
return [hashlib.sha256(c).hexdigest() for c in chunks]
random.seed(7)
original = bytes(random.getrandbits(8) for _ in range(200_000))
edited = b"HELLO" + original # insert 5 bytes at the front
for name, split in [("fixed", fixed_chunks), ("content-defined", cdc_chunks)]:
before = set(digests(split(original)))
after = digests(split(edited))
new = sum(1 for d in after if d not in before)
print(f"{name}: {len(after)} chunks after edit, {new} must be uploaded")
Output:
fixed: 49 chunks after edit, 49 must be uploaded
content-defined: 44 chunks after edit, 1 must be uploaded
With fixed chunks, a 5-byte insert forces all 49 chunks to be re-sent. With content-defined chunks, only the one chunk containing the insert changes.
| Fixed-size | Content-defined | |
|---|---|---|
| CPU cost | Very low | Higher (rolling hash per byte) |
| Insert or delete in the middle | Re-sends everything after the edit | Re-sends only nearby chunks |
| Dedup across similar files | Weak | Strong |
| Chunk sizes | Exactly equal | Vary between min and max |
| Good for | In-place overwrites, appends | Documents, archives, general files |
Deduplication by hash
Because a chunk's id is the SHA-256 of its contents, two users uploading the same video or the same installer produce the same hashes. The second upload is answered "already have it" and costs only a metadata row. SHA-256 collisions (two different chunks with the same hash) are so improbable that systems treat the hash as a unique id.
Deduplication brings two responsibilities:
- Reference counting and garbage collection. A chunk can be deleted only when no revision of any file uses it. Keep a reference count per block, or periodically mark all reachable blocks and sweep the rest. Deletions should be delayed (for example by days) so a race between "delete" and "a new commit references it" cannot lose data.
- Privacy side channel. Cross-user deduplication can reveal whether a file exists somewhere on the service: if an upload of a known file completes instantly, the uploader learns someone already has it. Mitigations include deduplicating only within one user or organisation, or always requiring the client to send the bytes for small or sensitive files. Proof-of-ownership schemes stop a client from claiming a chunk it only knows the hash of.
Common mistake
Trusting the client's hash. A malicious client could upload garbage under the hash of someone else's chunk, poisoning every file that uses it. The block service must recompute the hash of received bytes and reject mismatches.
Deep dive 2: the sync protocol and change notifications
Each device needs to learn about changes quickly without hammering the server.
Option 1: plain polling
Every device asks "anything new?" every 30 seconds. With 10 million connected devices, that is 10,000,000 ÷ 30 ≈ 333,000 requests per second, almost all answering "no". Latency is up to 30 seconds. Too wasteful at scale, but fine as a fallback.
Option 2: long polling
The device sends a request that the server holds open until either something changes or a timeout (say 60 seconds) expires. When a change arrives, the server answers immediately, and the client fetches the changes and opens a new long poll. Latency is near real time; idle cost is one request per device per minute, about 167,000 per second for 10 million devices, but each one is cheap and mostly waiting. Long polling works through almost every proxy and firewall because it is ordinary HTTP. Dropbox has publicly described using a long-poll notification endpoint for its desktop client.
Option 3: WebSocket or server-sent events
A WebSocket is a persistent two-way connection over a single TCP connection. The server pushes "namespace ns_42 changed" the instant it happens. This gives the lowest latency and the least request overhead, but needs connection-holding servers that track millions of sockets and a routing layer that knows which server holds which device.
Design the notification as a hint, not the data
The notification service only says "something changed in namespace X". The device then calls the change feed with its cursor. This separation has big benefits:
- Missed notifications are harmless. If a push is lost or the connection drops, the next feed call with the old cursor still returns every change. Correctness comes from the cursor, not from the push.
- The notification service can be stateless and lossy. It does not need to store anything durably.
- Batching: a hundred changes in one second produce one notification and one feed call.
device notify svc metadata svc
| subscribe(ns_42) | |
|------------------------>| |
| |<-- ns_42 changed --| (after commit)
|<-- "changes: true" -----| |
| GET /changes?cursor=c_1043 |
|--------------------------------------------->|
|<-- entries seq 1044..1051, cursor c_1051 ----|
| apply locally, save cursor |
Scaling the notification tier
Ten million idle connections need many servers, but each connection is cheap when idle. A pub/sub layer (for example Redis pub/sub, or a Kafka topic consumed by every notification server) carries "namespace changed" events. Each notification server subscribes to the namespaces of the devices connected to it. A shared folder with 500 members produces 500 wake-ups, which is fine; a folder shared with an entire 100,000-person company should be rate-limited or batched.
Deep dive 3: conflicts, versioning and offline edits
How a conflict happens
Alice and Bob both have plan.docx at revision 5. Alice's laptop is offline on a train. Both edit the file. Bob commits first and creates revision 6. When Alice reconnects, her client commits with parentRev: 5. The server sees that the latest is 6, not 5, and returns 409 Conflict.
rev 5 ---+---> rev 6 (Bob, committed first)
|
+---> Alice's edit (parentRev 5) -> CONFLICT
Resolving it
For arbitrary files, the server cannot merge two binary documents automatically. The standard, safe approach used by file sync products is the conflicted copy:
- Bob's revision 6 stays as
plan.docx. - Alice's version is committed as a new file, for example
plan (Alice's conflicted copy 2026-10-10).docx. - Both people see both files and decide manually.
No data is lost, and nothing is silently overwritten. "Last writer wins" by timestamp is simpler but silently destroys one person's work, and device clocks can be wrong, so it is a poor default for user documents.
Other strategies for special cases:
| Strategy | When it fits |
|---|---|
| Conflicted copy | General files, the safe default |
| Last writer wins | Unimportant metadata like "last opened" |
| Automatic merge | Formats the service understands, such as plain text with a three-way merge |
| Locking | Files that cannot be merged, for example CAD drawings; users check out a file before editing |
| CRDT or operational transform | Real-time collaborative editors (a different product) |
Folder-level conflicts
Some conflicts are structural: Alice renames /Docs to /Archive while Bob creates a file inside /Docs. Because files reference their parent folder by id, not by path, Bob's new file simply appears inside /Archive. Using ids for parents removes a whole class of conflicts. A harder case is Alice deleting a folder while Bob adds a file inside it; a safe rule is to keep the folder alive, or restore it, when it still contains a new file.
Versioning
Every commit creates a new row in revisions that lists its block hashes. Old revisions keep their block references, so restoring revision 16 is just a new commit whose block list equals revision 16's. Storage cost is modest because unchanged chunks are shared between revisions. A retention job deletes revisions older than 30 days (or per plan), decrements block reference counts, and lets garbage collection reclaim unreferenced chunks. Deleted files are just a revision with is_deleted, which makes "restore deleted file" easy.
Offline edits
The local index is what makes offline work. When the client starts or reconnects, it:
- Scans the sync folder and compares each file's size, modification time and, if needed, hash against the local index, to find local changes made while offline.
- Fetches remote changes since its cursor.
- For files changed only locally: commits them. Only remotely: downloads them. Both: conflict handling as above.
Order matters: the client must never overwrite a local unsynced edit with a remote version. It treats "local file differs from index" as precious data until it has been committed or saved as a conflicted copy.
Deep dive 4: sharing and permissions
Sharing a folder turns it into its own shared namespace with an access-control list (ACL): a list of users and their roles. Each member's client mounts the namespace at some point in their own tree, for example /Team Docs.
- Permission check: every metadata and block operation checks the caller's role on the namespace. Viewers can read but not commit. Cache ACLs in the metadata service with short expiry, because they are read on every request but change rarely.
- Block access: a block download must be authorised. One approach is to issue short-lived signed URLs only for blocks that belong to revisions the user can read; never allow "download any block if you know the hash", or hashes become secret keys.
- Link sharing: a public link maps a random, unguessable token to a file or folder and a role, with optional expiry and password.
- Revoking access: remove the ACL entry and close the user's subscriptions. Files already downloaded to their devices cannot be clawed back, except through device management in enterprise products.
Resumable uploads
Chunking already makes transfers resumable: the client uploads chunk by chunk and, after an interruption, calls blocks/check again to learn which chunks the server already has. The commit happens only when all blocks are present, so a half-uploaded file is never visible to other devices. For very large files, the client uploads several chunks in parallel and limits its bandwidth so it does not saturate the user's connection.
Scaling and bottlenecks
| Component | Pressure | Approach |
|---|---|---|
| Block storage | Hundreds of PB | Object storage with erasure coding; tiering for old versions |
| Block service | Upload bandwidth | Stateless, horizontally scaled, close to users; direct-to-storage signed URLs |
| Metadata DB | 200 billion rows, transactional commits | Shard by namespace; read replicas for listings |
| Huge namespaces | One shared folder with millions of files | Split hot namespaces, or paginate listings and journal reads |
| Notification tier | Millions of open connections | Many lightweight connection servers, pub/sub fan-out |
| Change feed | Many devices catching up | Range scans by (namespace_id, seq); compacted snapshots for new devices |
| Hashing on the client | CPU on laptops | Hash only changed files; do it in the background |
New device bootstrap. A brand new device does not replay the journal from seq 1. It downloads a snapshot of the current tree with the current cursor, then follows the journal from there.
Failure handling
- Upload interrupted: resume by checking missing blocks; nothing is committed until all blocks exist.
- Client crash mid-commit: the commit is one database transaction, so it either happened or not. The client re-checks the latest revision on restart.
- Duplicate commit (client retried after a timeout): include an idempotency key or rely on
parentRev: the second identical commit either is recognised as a duplicate or gets a conflict that the client can resolve by seeing the server already has its content. - Notification lost: harmless; the client also does a periodic catch-up call with its cursor.
- Metadata shard down: that shard's namespaces are temporarily read-only or unavailable; other users are unaffected. Each shard is replicated with automatic failover.
- Block storage corruption: replicas or erasure-coded pieces rebuild it; hashes detect corruption on every read.
- Garbage collector bug: the most dangerous failure, because it deletes real data. Use delayed deletion, soft deletes and audits that compare reference counts with a full mark phase before anything is purged.
Trade-offs and alternatives
| Decision | Option A | Option B | Choice and reason |
|---|---|---|---|
| Chunking | Fixed-size | Content-defined | CDC for general files; fixed for in-place formats |
| Chunk size | Small (1 MB) | Large (8 MB) | Around 4 MB: fewer metadata rows than 1 MB, better deltas than 8 MB |
| Dedup scope | Global | Per user or organisation | Per organisation for privacy, global only for public content |
| Notifications | Polling | Long poll or WebSocket | Long poll or WebSocket, with polling as fallback |
| Conflict policy | Last writer wins | Conflicted copy | Conflicted copy: never lose user data |
| Metadata DB | Sharded SQL | NoSQL | Sharded SQL for transactions within a namespace |
| Mobile sync | Download everything | Download on open | On demand, to save phone storage and data |
| Encryption | Server-side keys | End-to-end | Server-side keeps dedup, search and previews; end-to-end protects more but loses them |
What interviewers probe
"Why not just store each file as one object?" You lose delta sync, deduplication and fine-grained resume. Every small edit to a large file would re-upload the whole thing.
"How do you know which device made a change, so it does not download its own upload?" Each commit records the author device. When that device reads the journal, it skips entries it authored, or simply notices it already has the matching revision in its local index.
"What if a user has a million files?" Listing must paginate. The journal still works because devices read it in pages by seq. The initial snapshot is streamed in pages too.
"How do you enforce quotas with deduplication?" Charge each user for the logical size of their files, not for the unique bytes they cause to be stored. Otherwise quota depends on what other people uploaded, which is confusing and leaks information.
"How would you search file contents?" An asynchronous worker reads new revisions, extracts text, and indexes it in a search engine, partitioned by namespace and filtered by ACL at query time.
"How do you sync over a slow, metered connection?" Pause large uploads on metered networks, prioritise small files and metadata, compress chunks before upload, and let the user choose which folders sync to which device (selective sync).
Interview questions
Q1. Why separate the metadata service from block storage?
They have opposite needs. Blocks are large, immutable and need cheap, durable bulk storage, which object storage provides. Metadata is small, changes constantly and needs transactions and strong ordering, which a sharded database provides. Separating them lets each scale and fail independently.
Q2. What is content-defined chunking and why is it better than fixed-size?
It places chunk boundaries where a rolling hash of nearby bytes matches a pattern, rather than every N bytes. An insertion or deletion only changes the chunks around the edit, because later boundaries move with the content. With fixed-size chunks, a one-byte insert shifts every later chunk and forces a full re-upload.
Q3. How does deduplication work?
Each chunk is identified by the cryptographic hash of its contents. Before uploading, the client asks which hashes the server lacks and sends only those. Identical chunks across versions, files and possibly users are stored once, with reference counts so they are deleted only when unused.
Q4. How do other devices learn about a change?
They keep a long-poll or WebSocket connection to a notification service, which sends a lightweight hint that a namespace changed. The device then reads the change journal from its last cursor. Because the cursor guarantees completeness, a lost notification only delays sync; it never loses changes.
Q5. What is a cursor in the sync protocol?
It is the position in a namespace's ordered change journal that a device has fully applied. Asking for changes "since cursor 1043" returns every change after it, in order. It makes sync restartable and idempotent.
Q6. How do you detect a conflict?
Every commit includes the revision the client started from. If the server's latest revision differs, someone else committed in between, and the server rejects the commit with a conflict. This is optimistic concurrency control: no locks, just a version check at commit time.
Q7. How do you resolve a conflict?
For general files, keep the first committed version under the original name and save the other as a "conflicted copy" next to it, so no work is lost. Merging is only safe for formats the system understands, and last-writer-wins is reserved for unimportant data.
Q8. How does version history avoid storing full copies?
Each revision is just an ordered list of chunk hashes. Unchanged chunks are shared between revisions, so a new version costs only its changed chunks plus a metadata row. Restoring an old version means committing its block list again.
Q9. How do offline edits work?
The client keeps a local index of what it last synced. On reconnect, it compares disk contents with the index to find local edits, fetches remote changes since its cursor, and commits, downloads or creates conflicted copies as appropriate. It must never overwrite an unsynced local change.
Q10. How do you safely delete chunks with deduplication?
Track references, either with counts or with a periodic mark-and-sweep over all live revisions. Delete only chunks with no references, and only after a delay, so a race with a new commit that references the same chunk cannot destroy data.
Q11. What chunk size would you choose?
Around 4 MB is a common middle ground. Smaller chunks give finer deltas and better dedup but multiply metadata rows and requests; larger chunks reduce overhead but re-upload more for small edits. Content-defined chunking typically targets an average size with minimum and maximum limits.
Q12. Is cross-user deduplication safe?
It saves storage, but it can reveal whether a particular file already exists on the service, and a client that knows a hash might claim a chunk it does not own. Many designs limit dedup to one user or organisation and require proof of ownership.
Q13. How do you shard the metadata?
By namespace, so all files, revisions, journal entries and ACLs of one sync root live on one shard. Commits then stay single-shard transactions. Very large shared namespaces may need special handling, such as dedicated shards.
Key takeaways
- Split the system into a strongly consistent metadata service and an immutable, content-addressed block store.
- Chunk files and identify chunks by SHA-256; upload only missing chunks, then commit a new block list atomically.
- Content-defined chunking survives insertions; fixed-size chunking does not.
- Each namespace has an ordered change journal; devices sync with a cursor, and notifications are only hints.
- Use optimistic concurrency (
parentRev) to detect conflicts, and conflicted copies to resolve them without data loss. - Versions are cheap because revisions share chunks; garbage collection must be careful and delayed.
- Shard metadata by namespace; shared folders are their own namespaces with ACLs.
- Verify hashes on the server and authorise every block read.
Next lesson
Continue with the course overview.

