AI System Design
← All chapters
Chapter 15 7 min read

Design Google Drive

A cloud file store must keep every user's files durable, private and identical across phones, laptops and the web, while spending as little bandwidth and disk as possible. The core ideas are splitting files into content-addressed blocks so only changed pieces move, keeping a strongly consistent metadata database as the source of truth, and pushing change notifications to every device so they can pull what is new.

Architecture at a glance
  1. Client app (watcher + chunker)
  2. Load balancer
  3. Block servers (split, compress, encrypt)
  4. Cloud object storage (S3, cross-region)
  5. API servers
  6. Metadata DB + cache
  7. Notification service (long polling)
  8. Other devices pull changes

Problem & requirements

The product lets a user add a file once and see it everywhere: upload and download, automatic sync across devices, a history of revisions, sharing with friends or colleagues, and a notification when a shared file is edited, deleted or shared. Real-time collaborative editing in the style of Google Docs is explicitly out of scope; this is about whole files that change occasionally, from small documents up to perhaps 10 GB.

The non-functional bar is unusually strict. Data loss is unacceptable, because users treat the service as the copy of record. Sync should feel fast, otherwise people stop trusting it. Bandwidth matters because many clients are on metered mobile links. And the system must scale to tens of millions of daily users and very large aggregate storage without the cost growing faster than the user base.

  • Functional: upload, download, sync, revisions, share, notify on change.
  • Non-functional: reliability and durability, fast sync, low bandwidth, scalability, high availability.

Back-of-the-envelope

Assume 50 million signed-up users and 10 million daily actives, each given 10 GB of free space. The theoretical ceiling is 50M x 10 GB = 500 PB, which tells us immediately that storage cost, not compute, dominates the budget and that deduplication and cold tiers are first-class concerns rather than optimisations for later.

If each active user uploads two files a day of about 500 KB, that is 10M x 2 = 20M uploads per day. Dividing by 86,400 seconds gives roughly 230 uploads per second, and assuming peaks are twice the average, about 460 per second. Double-checking: 20,000,000 / 86,400 is a little over 231, so the order of magnitude holds. Write traffic is modest; the challenge is correctness and storage efficiency, not raw request rate.

APIs and evolving from one server

Three API families cover the product. Uploads come in two flavours: a simple upload for small files sent in one request, and a resumable upload for large files, where the client first asks for a session URL, then streams data and can resume from the last acknowledged byte after a network drop. Downloads take a path or file ID. A revisions endpoint lists past versions with a limit so a heavily edited file does not return thousands of rows. Every call is authenticated and runs over HTTPS.

A single machine with a web server, a relational database and a local disk is a fine starting point, but it runs out of disk within weeks and is a single point of failure. The natural evolution is to move file bytes to object storage such as Amazon S3 with cross-region replication, so a whole-region outage does not lose data; put a load balancer in front of stateless web servers; move metadata to a replicated database; and keep the two stores separate so each can scale on its own axis.

Figure 1Two paths: blocks and metadata
Two paths: blocks and metadatachangedblocksPUT byhashlifecyclecommit metadatablocks storedread /writeon missfile changedpollreturnsifofflinepull new metadataClient Afile watcher + chunkerBlock serverschunk · compress · encryptCloud storageblocks keyed by hashCold storageold, rarely read versionsAPI serversauth, metadata, sharingMetadata cacheMetadata DBfiles · blocks · versionsClient Bsame user, other deviceNotification servicelong pollingOffline backup queueevents for absent devices
File bytes go through block servers into object storage, while the small metadata commit goes through API servers to a strongly consistent database. Only after the commit does the notification service wake other devices, which then pull just the blocks they lack.

Deep dive: block servers, delta sync and dedup

Uploading a 2 GB video again because one byte changed would be wasteful, so files are cut into blocks of up to about 4 MB, each identified by a hash of its content. Block servers compress each block (the algorithm depends on file type), encrypt it, and store it in object storage under that hash. A file is then simply an ordered list of block hashes, and a new version is a new list.

This representation gives two big wins. Delta sync: when a file changes, the client recomputes block hashes and only uploads blocks whose hashes are new, so a small edit costs one block rather than the whole file. Deduplication: if two blocks anywhere in an account (or, with care, across accounts) share a hash, they are stored once. The split work can happen either on the client, saving upload bandwidth, or on block servers, which keeps clients simpler; most real systems push chunking to the client for exactly the bandwidth reason.

  • Fixed-size blocks are simple; content-defined chunking (rolling hash) survives insertions that shift every later byte.
  • Encrypt after compressing — encrypted data looks random and will not compress.
Figure 2Delta sync of one edited file
Delta sync of one edited fileBlock 14 MBhash matches · skipBlock 24 MBedited · uploadBlock 34 MBhash matches · skipBlock 42 MBappended · uploadtotal 14 MB
The client hashes each fixed-size block and compares against the block list in metadata; here only blocks 2 and 4 changed, so 6 MB crosses the network instead of 14 MB. Identical hashes anywhere in the account are stored once (dedup).

Deep dive: metadata database and consistency

Metadata is the system's source of truth: which blocks make up which version of which file in whose namespace. If two devices could read different answers, a user might open a stale version or, worse, a sync client might delete a file it wrongly believes was removed. So the metadata store must be strongly consistent, and caches in front of it must be invalidated on write. A relational database gives ACID transactions for free, which is why it is the default choice here over an eventually consistent NoSQL store.

A workable schema has a user table; a device table holding a push ID per device; a namespace table representing a user's root directory; a file table with the latest metadata; a file_version table whose rows are immutable history; and a block table mapping each version to an ordered list of block hashes. Making version rows read-only protects revision history from accidental corruption by later writes.

Upload, download and notification flows

An upload runs two requests in parallel. The first adds file metadata with status pending and tells the notification service a change is coming. The second sends the content to block servers, which chunk, compress, encrypt and write blocks to cloud storage. When storage confirms, a callback flips the status to uploaded and notifies every other device of the owner and of anyone the file is shared with.

Downloads are triggered by those notifications. An online client receives a change event and asks API servers for the new metadata, then fetches only the blocks it does not already have and reassembles the file. Notifications use long polling rather than WebSocket because traffic is one-way and infrequent: each client holds a request open, the server answers when something changes or the timeout fires, and the client immediately reconnects. Changes for offline clients go into an offline backup queue so they are delivered on reconnect.

Figure 3Upload, then sync another device
Upload, then sync another deviceClient AAPI serverBlock serversCloud storageNotificationClient B1. add file metadata (status: pending)2. upload changed blocks3. compress + encrypt each block4. PUT blocks keyed by content hash5. all blocks stored6. commit version 7 (status: uploaded)7. file 42 changed8. long poll returns: changes ready9. fetch version 7 block list10. download missing blocks only11. decrypted, decompressed blocks
Metadata is written twice: as pending before any bytes move and as uploaded once every block is safe, so a crash mid-upload never exposes a half-written file. Client B downloads only the blocks its local copy is missing.

Saving storage and resolving conflicts

With hundreds of petabytes in play, every percent of storage saved is real money. Deduplicate identical blocks at the account level. Cap the number of revisions kept and give recent versions more weight, because a frequently edited file could otherwise accumulate thousands of near-identical versions. Move data that has not been touched for months to a cold storage tier such as S3 Glacier, which is far cheaper per gigabyte but slower to read.

When two users or two devices edit the same file concurrently, the rule is simple: the first version processed wins, and the later one is preserved as a conflict copy alongside the original. The losing user sees both and merges by hand. Automatic merging of arbitrary binary formats is impossible, and silently discarding work would break the durability promise, so presenting both copies is the honest trade-off.

Failure handling & wrap-up

Each component needs a story for its failure. A load balancer pair uses heartbeats so a secondary takes over. Block and API servers are stateless, so others pick up the traffic, although an in-flight upload may need to resume. Object storage is replicated across regions. The metadata database uses primary-replica replication with promotion of a replica if the primary dies. Notification servers hold many long-lived connections; when one fails, clients reconnect to another, which is slow if millions do so at once, so reconnection must be spread out.

The interview takeaway is the separation of concerns: bytes go to cheap, durable object storage in content-addressed blocks; meaning lives in a small, strongly consistent metadata database; and freshness is delivered by a lightweight push channel. A good extension to discuss is uploading directly from the client to object storage with signed URLs, which removes block servers from the data path but moves chunking, compression and encryption logic onto every client platform.

Key numbers

Block size
~4 MB
Ceiling storage (50M users x 10 GB)
500 PB
Upload QPS (avg / peak)
~230 / ~460
Max file size supported
~10 GB
Daily active users
10 million

Key terms

Block server
A service that splits files into blocks, compresses and encrypts them, and stores them in object storage keyed by content hash.
Delta sync
Uploading only the blocks whose hashes changed between versions instead of the whole file.
Deduplication
Storing a block once when multiple files or versions contain identical content, detected by equal hashes.
Resumable upload
An upload protocol that lets a client continue from the last acknowledged byte after an interruption.
Long polling
A client holds an HTTP request open until the server has news or a timeout fires, then immediately reconnects.
Namespace
The root directory that scopes a user's files, used to partition and authorise metadata.
Conflict copy
A preserved second version created when a concurrent edit loses to the first-processed version.
Cold storage
A cheap, slow storage tier for data that is rarely accessed.

Common mistakes

  • Using an eventually consistent store for metadata, which lets devices disagree about which version is current.
  • Re-uploading whole files on every edit instead of syncing changed blocks.
  • Encrypting before compressing, which destroys the compression ratio.
  • Choosing WebSocket for notifications without justifying it; traffic is one-way and sparse, so long polling suffices.
  • Forgetting offline devices — changes must be queued and replayed when they reconnect.
  • Keeping unlimited revisions, which silently multiplies storage cost.

Further study

  • Dropbox: How We've Scaled Dropbox (Kevin Modzelewski talk)
  • Dropbox Magic Pocket exabyte storage system
  • rsync algorithm (Andrew Tridgell)
  • Amazon S3 Cross-Region Replication
  • LBFS: A Low-bandwidth Network File System (2001)

Now practise it