LearnHLDDesign a file sync service

Design a file sync service

You edit a spreadsheet on a laptop with the wifi off. Someone else edits the same file from their desk. You land, open the lid, and the two versions have to become one thing without anybody being asked to sort it out.

Uploading a file is the easy part of this question. Keeping one folder identical across four devices, two of which are asleep, one of which is on hotel wifi, and one of which has been switched off for a month, is the part that gets designed badly.

Step 1: Understand the problem

Three of the six answers below decide the whole design, and the first one decides the bandwidth bill.

You askThey sayWhat it settles
How big are files, and how much of one changes in an edit?A 2KB note up to a 50GB video, and an edit usually touches a few percent.Chunking. You never upload a whole file after a one line change, so files are split into blocks and only changed blocks move. That single answer removes most of the traffic in the system.
How many devices per user, and are they online?Three or four, and most of them are asleep most of the time.The server cannot push to a device that is not there, so syncing is a device pulling a change log from a cursor it holds. A nudge wakes it up, the nudge carries no data.
What happens when two devices edit the same file?Keep both. Never silently lose one.Conflicted copies rather than a merge. Say this out loud: a general file sync service cannot merge a binary file, so the honest answer is to keep both and name them clearly.
Do we deduplicate the same bytes across different users?Yes. Storage is the second largest line on the bill.Content addressed blocks, and a privacy problem you should name: if an upload is instant because someone else already has that block, you have told the uploader that file exists on the service.
How fast must a change reach another device?A few seconds when both are online.A long lived notification channel separate from the data path. The nudge is tiny and urgent, the transfer is large and allowed to be slow, and mixing them is how one slow upload delays everybody.
Do we need version history?30 days, plus undelete.Almost nothing extra, which is the quiet payoff of chunking. Blocks are immutable and content addressed, so a version is a list of hashes and a delete is a flag, not a destructive operation.

What you are building, and what you cut

In scope
  • Keep a folder identical across a user's devices. Including devices that have been offline for weeks.
  • Upload and download with resume and dedupe. Nobody restarts a 50GB upload because a train went into a tunnel.
  • Handle concurrent edits without losing data. A conflicted copy, named so a human can tell which is which.
  • Share a folder with other users. Which is where the permission model and the fan out both get interesting.
The numbers you commit to
  • 500 million users, a few hundred files each.
  • A change on one device appears on another in under 5 seconds when both are online.
  • A device offline for a month catches up without downloading the folder again.
  • No file is ever lost, including when two people edit at once.
Cut, and say so out loud
  • Collaborative editing of the file contents. Different problem, different algorithm, and it has its own page.
  • Full text search inside documents. An indexing pipeline that reads this system rather than part of it.
  • Fine grained per file permissions. Folder level sharing is enough to show the design.
  • Selective sync policy on mobile. A client product decision with no interesting server side.

Back of the envelope

What actually moves, and what actually piles up
500 million
500
12 MB
Files stored500M x 500 = 250 billion
Logical bytes250B x 12MB = 3.0 EB
After cross user dedupe, about 25%2.3 EB
Metadata rows at 3 versions each250B x 3 = 750 billion
Metadata size at ~300 bytes a row225 TB
750 billion rows
of metadata, which is a harder problem than the exabytes of file content

Two things fall out of this. The bytes are object storage, which is a solved problem you can buy by the petabyte, while 750 billion small rows that have to be read in order, per device, from a cursor, is the part that gets redesigned twice. And chunking is the entire bandwidth story: a one line change to a 12MB document moves one 4MB block, and a one line change to a 2GB video also moves one 4MB block.

The number people get wrong

Candidates size the block store, find an impressive number of petabytes, and design around it. Blocks are the easy half: they are immutable, content addressed, and you can buy that storage. The metadata is the half that hurts, because every device asks “what changed since I last asked” and that question has to be answered in order, cheaply, billions of times a day, for namespaces that range from one file to two million.

Step 2: Propose the high level design

The API

Five calls, and the first one is the one most designs are missing.

POST/v1/blocks:check
{
  "hashes": ["a3f1...", "7c90...", "1b2e..."]
}
returns 200 { "missing": ["7c90..."] }
Why: The client asks before it uploads. On a corporate network where the same 2GB installer sits in forty people's folders, the thirty ninth person uploads nothing at all. This one call is the whole deduplication win and it costs a hash computation the client was doing anyway.
PUT/v1/blocks/{sha256}
returns 201
Why: The name is the hash of the content, so the write is idempotent by construction. Uploading the same block twice is free, a retry after a timeout is safe, and a corrupted transfer fails its own checksum rather than silently storing garbage.
POST/v1/commit
{
  "path": "/work/report.xlsx",
  "blocks": ["a3f1...", "7c90..."],
  "baseVersion": 41207
}
returns 200 { "version": 41208 } or 409 { "serverVersion": 41209 }
Why: Metadata is committed last, after every block is durable, so a file never appears in the namespace pointing at bytes that are not there yet. The baseVersion is how a concurrent edit is detected, and that 409 is what becomes a conflicted copy.
GET/v1/changes?cursor=41207
returns 200 { changes: [...], cursor: 41260, hasMore: false }
Why: A device pulls an ordered log from the cursor it holds. A laptop that was off for a month asks once with an old cursor and receives everything since, in order, with no special catch up path and no tree comparison. This is the single most important endpoint in the system.
GET/v1/longpoll?cursor=41207
returns 200 { "changed": true } after up to 90 seconds
Why: A cheap, long lived connection whose only job is to say "something happened, come and ask". Keeping the nudge separate from the data keeps the push tier tiny, stateless and easy to restart, which matters when you are holding a connection open for every awake device on the planet.

The data model

Two stores, split by what changes rather than by what the data is.

namespace_journalsharded by namespace, append only, cursor ordered
namespace_idbigintPKPartition key. One user's private folder or one shared folder. Everything about it lives together, which is what makes a cursor read a single partition scan.
seqbigintPKClustering key, strictly increasing within the namespace. This number is the cursor every device remembers.
pathvarchar(1024)IDXThe current path. A rename is an entry, not an update, which is why history survives one.
block_listarray of sha256The file is this list. The bytes live somewhere else entirely and are shared with whoever else has them.
size, mtimebigint, timestampEnough for a client to decide what to fetch without opening anything.
deletedbooleanA delete is an append like every other change. That is what makes undelete and 30 day history fall out for free.
Sample row
9912 | 41208 | /work/report.xlsx | [a3f1..., 7c90...] | 184320 | false
Nothing in this table is ever updated in place. A device syncs by reading forward from its cursor and applying what it finds, which means the hardest client bug in this domain, two devices disagreeing about what the truth is, cannot happen: the journal is the truth and the cursor says how much of it you have seen.
blocksobject storage, keyed by content hash
sha256char(64)PKThe name is the content. Two users with the same holiday photo store it once.
sizeintUp to 4MB. Bigger blocks mean less metadata and worse delta efficiency, and this number is the knob.
refsapproximateNot a transactional reference count. Counting references to a block from a billion namespaces exactly is a distributed counter nobody wants.
Sample row
a3f19c... | 4194304 | many
Deleting a block the moment nothing points at it is a distributed reference count, and getting it wrong deletes somebody's data. So nothing is deleted synchronously: a mark and sweep job runs offline, and blocks are kept for longer than the longest version history. Storage is cheaper than the incident.

The whole system on one whiteboard

Figure 1. Two paths that meet only at the commit. Blocks go up the top of the picture and are allowed to be slow. The journal entry goes across the middle and is what makes the change real for every other device.

Walking Figure 1:

  1. The client notices a file changed and splits it into blocks. This happens on the device, before anything touches the network, and how it splits is the first deep dive.
  2. It asks which of those block hashes the server is missing. Usually most of them are already there, either from the previous version of this file or from another user entirely.
  3. Only the missing blocks are uploaded, straight to the block store, resumable and in parallel.
  4. Then, and only then, the client commits: a path, an ordered list of hashes, and the version it believed it was editing.
  5. The metadata service checks that version against the journal.
  6. If it matches, the change is appended. The file is now real, and it was never real at any moment when its bytes were missing.
  7. The notifier learns a namespace moved.
  8. Every awake device holding a long poll for that namespace is told, with no payload.
  9. Each of them asks for changes since its own cursor, and gets exactly what it missed.

The three dashed arrows are the background: membership decides who sees a namespace, history is a view over old journal entries, and a sweeper eventually deletes blocks nothing points at any more.

Your answer

Step 2 asks the server which blocks it already has before uploading anything. What does that one question buy you, and what does it quietly leak?

Step 3: Design deep dive

Why fixed size blocks are the wrong answer

Split a file every 4MB and you get a design that works beautifully until someone inserts a line at the top of a document. Every byte after the insertion shifts, so every block boundary lands in a different place, and a one character edit uploads the entire file.

The fix is to let the content decide where the boundaries go.

1/4 A small window slides over the file one byte at a time. The hash of the window is cheap to update as it moves, which is what makes this affordable on a 50GB video: one add and one subtract per byte rather than a fresh hash.
Figure 2. Boundaries are chosen by what the bytes are, not by how far through the file you are. Insert something at the top and exactly one chunk changes, because the boundaries after the insertion land in the same places they did before.

Be honest about the cost. Variable chunks mean you cannot compute where a byte is without walking the list, and the rolling hash is real CPU on the client. Both are fine. The alternative is uploading a 2GB video because somebody renamed a layer in it.

The follow up you will get

“What chunk size would you pick?” Smaller chunks deduplicate better and move less on an edit, and cost you more metadata rows and more round trips per file. Say the shape of the trade and give a range rather than a number: somewhere around 1MB to 8MB average, tuned by measuring the ratio of bytes moved to bytes changed on real files. Then name the detail that shows you have thought about it, which is that you also need a minimum and maximum chunk size, because a pathological file can produce a boundary every 50 bytes or none at all for a gigabyte.

Syncing by cursor, not by comparison

The tempting model is that a client compares its folder to the server’s folder and fixes the differences. That is correct, and it costs you a full comparison on every sync, needs the client to hold the entire tree, and gets slower as the folder grows even when nothing changed.

The journal replaces it with a number. Every device remembers one integer per namespace, and syncing means asking what came after it.

1/6 The phone is awake and holding a connection open that costs almost nothing. It is not asking for data, it is asking to be told when there is some.
Figure 3. A laptop that was closed for a month and a phone that was awake the whole time run exactly the same code path. One of them just has a smaller number.

The property worth naming out loud: there is no catch up path. Offline for an hour and offline for a year run the same code with a different number in it. The rarely exercised path and the constantly exercised path are the same path, which is how you avoid the bug that only shows up after somebody’s long holiday.

Two devices, one file

The 409 from the commit endpoint is where a product decision lives, and the interviewer is asking which one you make.

Concurrent edit policy
The loser writes its version to a new path, "report (conflicted copy from Sam's laptop).xlsx", and both appear on every device. It looks unsophisticated and it is the right answer: no data is lost, the user can see exactly what happened, and a human resolves it with the context a server does not have.

Pick conflicted copies, and say the sentence that justifies it: the server knows the bytes, the user knows the intent, and only one of those is qualified to throw work away.

Break it

Sync pressure
normal
normalshared team foldersomeone pastes 200k filesa client goes into a looprepaired
Healthy. Most devices are asleep. The awake ones hold a long poll that does nothing, commits are a few rows, and the block store is handling large but entirely unexciting transfers.

Trade-offs

ChoiceWhat you gainWhat you payPick it when
Content defined chunkingAn insert at the top of a file changes one chunk instead of all of them, and identical regions dedupe across users.Rolling hash CPU on the client, variable chunk sizes, and no way to seek to a byte offset without walking the list.Any sync product where files are edited rather than only added. For write once storage, fixed blocks are simpler and fine.
Cursor journal over tree comparisonSync cost is proportional to what changed, not to how many files exist, and offline for an hour and offline for a year are the same code path.A journal per namespace is a hot partition for bulk operations, and the log has to be compacted eventually.Always. Tree comparison looks simpler and gets slower every month the product succeeds.
Conflicted copies over mergingNo edit is ever lost, and the behaviour is explainable to a user in one sentence.Users occasionally see two files and have to decide, which feels like the product failing at its job.Whenever the service does not understand the file format, which for a general sync service is always.
Cross user deduplicationA large cut in stored bytes, and instant uploads for files the service has seen before.An upload that completes suspiciously fast tells the uploader that someone else has that exact file.Consumer storage, usually yes. If that leak matters, scope dedupe to one user and accept the bill.

Interview replay

Interviewer
Someone edits a file on their laptop. Take me through what reaches the server.
Opening. They want to know whether you upload the file.
You
Not the file. The client chunks it with a rolling hash, so boundaries are chosen by content, hashes each chunk, and asks the server which of those hashes are missing. Usually almost none are, because most of the file did not change. It uploads the missing chunks to the block store, and only after they are durable does it commit metadata: the path, the ordered list of hashes, and the version it thought it was editing.
Blocks before metadata, stated as an ordering with a reason. That ordering is the whole correctness story.
Interviewer
Why not split the file every four megabytes?
Checking whether content defined chunking was a decision or a phrase.
You
Because an insert at the start shifts every byte, so every fixed boundary lands somewhere new and the whole file looks changed. With a rolling hash the boundaries are decided by the bytes around them, so after the inserted region the window sees the same content it saw before and re-synchronises. One chunk changes instead of three thousand.
Explains the re-synchronisation, which is the part people repeat without understanding.
Interviewer
How does another device find out?
The sync model question.
You
It holds a long poll that returns a bare "something changed", then asks for changes since its own cursor. I would keep the payload out of the notification deliberately, because the moment the push carries data it has to know what each device has already seen, and then a device that was asleep for twenty changes needs a different code path. With a cursor, a device offline for an hour and one offline for a year do exactly the same thing.
Refusing to put data in the push, and explaining it as path reduction rather than as purity, is the senior framing.
Interviewer
Two devices edit the same file while both are offline. Then both come back.
The conflict question. There is a right answer and it is not clever.
You
The first commit wins, the second gets a 409 because its base version is stale, and the client writes its version to a conflicted copy with the device name in the filename. I would not try to merge. A sync service does not know whether it is holding a text file or a video project, and a bad automatic merge destroys work in a way that is much harder to recover from than two files sitting next to each other.
Picks the unglamorous answer and defends it on blast radius rather than on difficulty.
Interviewer
Fifty people share a folder and someone drops two hundred thousand files into it.
The hot partition question, in its file sync costume.
You
The namespace is one partition on purpose, so that is 200,000 appends to one shard, and then fifty devices waking up to read it. I would rate limit writes per namespace, batch a bulk import into fewer journal entries rather than one per file, and jitter the wake ups so fifty clients do not arrive in the same 200 milliseconds. What I would not do is split the namespace across shards, because the ordered cursor is the thing that makes the whole client simple.
Names the mitigations and then names what they refuse to give up, which is the harder half of the answer.
Interviewer
How long would you keep deleted blocks before collecting them?
Open ended. Checking for an invented number.
You
Longer than the version history window, and I would want to measure rather than assert. The risk is asymmetric: collecting too late costs storage, which is cheap, and collecting too early loses data, which is unrecoverable and ends up in a news article. So I would start at something like twice the retention window, run the sweeper in a dry run mode that only reports what it would delete, and compare that against expectations for a while before letting it actually delete anything.
The dry run is what someone who has written a destructive background job says. It costs one sentence and it is very hard to fake.

Checkpoint

Checkpoint

1. Why are chunk boundaries chosen by a rolling hash rather than by byte offset?

2. Why must blocks be uploaded before the metadata commit, rather than after?

3. Same design, but it now syncs a single multi gigabyte database file that changes constantly. What breaks first?

Worth memorising
  • 4MB blocks: a one line edit to a 2GB file moves 4MB. That is the whole bandwidth argument.
  • Hundreds of billions of metadata rows, which is a harder problem than the exabytes of content.
  • One cursor per device per namespace. Offline for an hour and offline for a year are the same code path.
  • Blocks first, metadata last. A file must never exist in the namespace pointing at bytes that are not there.
Say this in 60 seconds

A file sync service is a chunking problem and a cursor problem, and almost nothing else. The client splits each file using a rolling hash so boundaries are decided by content, which means inserting a line changes one chunk instead of every chunk after it. It asks the server which chunk hashes are missing, uploads only those, and commits metadata last, so a file never exists in the namespace pointing at bytes that are not there. Metadata is an append only journal per namespace, and every device remembers one cursor into it, so syncing is asking what came after that number. A device offline for an hour and one offline for a year run exactly the same path. Notifications are a separate long poll that carries no data, just a nudge, which keeps the push tier stateless. Concurrent edits produce a conflicted copy rather than a merge, because a sync service does not know the file format and losing work silently is much worse than showing two files. The failure I would call out is that a namespace is deliberately one partition, so a shared folder with a bulk import is a hot shard, and I would rate limit and batch rather than give up the ordered cursor.

IndGeek provides solutions in the software field, and is a hub for ultimate Tech Knowledge.