Google Drive is one of the most comprehensive system design problems — it covers file storage, synchronization, conflict resolution, sharing, deduplication, and access control. Let me walk through how I'd approach designing it.
Functional Requirements
- Upload and download files
- Sync files across multiple devices
- Share files with other users (read/write permissions)
- View revision history and restore previous versions
- Files up to 5GB
Non-Functional Requirements
- 1 billion users, 100 million DAU
- Average storage per user: 10 GB → total 10 exabytes
- Upload/download with high throughput
- Strong consistency within a device session, eventual consistency across devices
- 99.99% availability
Core Architecture
Client (Web/Desktop/Mobile)
↓
API Gateway + Load Balancer
↓
├── Metadata Service → MySQL (file tree, permissions, versions)
├── Upload Service → Blob Store (S3, GCS)
├── Sync Service → Kafka + Redis
└── Notification Service → WebSocket / SSE
File Chunking
Uploading a 5GB file as a single unit is fragile — any network interruption means restarting from scratch. Instead, we chunk files into fixed-size blocks (e.g., 4MB):
5GB file → 1,280 chunks of 4MB each
Each chunk uploaded independently
Progress tracked per chunk
Resume from last successful chunk on failure
At the client:
def upload_file(file_path):
chunks = split_into_chunks(file_path, chunk_size=4*1024*1024)
for i, chunk in enumerate(chunks):
chunk_hash = sha256(chunk)
# Ask server if this chunk already exists (deduplication)
if not server.has_chunk(chunk_hash):
server.upload_chunk(chunk_hash, chunk)
server.register_chunk(file_id, i, chunk_hash)
server.finalize_upload(file_id)Deduplication
If two users upload the same file (or two chunks are identical), we store only one copy. This dramatically reduces storage costs.
Block-level deduplication: Each chunk is identified by its SHA-256 hash. Before uploading, the client asks the server if the hash already exists. If yes, no upload needed — just reference the existing block.
chunk_store table:
| chunk_hash (SHA256) | storage_path | ref_count | size |
file_chunks table:
| file_id | chunk_index | chunk_hash |
When ref_count drops to 0 (all files referencing this chunk are deleted), the block is eligible for garbage collection.
File-level deduplication: If the full file hash matches an existing file, the new file just references the same chunks. The storage cost is O(1) regardless of file size.
Metadata Service
The metadata layer stores everything except the actual file bytes:
CREATE TABLE files (
file_id UUID PRIMARY KEY,
user_id UUID NOT NULL,
parent_id UUID, -- NULL for root
name VARCHAR(255),
type ENUM('file','folder'),
size_bytes BIGINT,
created_at TIMESTAMP,
updated_at TIMESTAMP,
is_deleted BOOLEAN DEFAULT FALSE
);
CREATE TABLE file_versions (
version_id UUID PRIMARY KEY,
file_id UUID,
version_num INT,
chunk_list JSON, -- ordered list of chunk_hashes
created_at TIMESTAMP,
created_by UUID
);
CREATE TABLE file_shares (
share_id UUID PRIMARY KEY,
file_id UUID,
shared_with UUID,
permission ENUM('viewer','editor'),
expires_at TIMESTAMP
);Use MySQL with master-replica replication. The file tree is a recursive structure — queries like "get all files in folder" need careful indexing:
-- Path enumeration or adjacency list for folder hierarchy
CREATE INDEX idx_files_parent ON files(parent_id, user_id, is_deleted);Upload Flow
Resumable upload protocol (similar to Google's):
- Client initiates upload:
POST /upload/initiate→ returnsupload_id - Client queries which chunks are missing:
GET /upload/{upload_id}/status - Client uploads missing chunks in parallel:
PUT /upload/{upload_id}/chunk/{index} - Client finalizes:
POST /upload/{upload_id}/complete
This handles network failures gracefully — resume by skipping chunks the server already has.
Parallel chunk uploads:
with ThreadPoolExecutor(max_workers=8) as executor:
futures = [executor.submit(upload_chunk, upload_id, i, chunk)
for i, chunk in enumerate(chunks)
if not server.has_chunk(upload_id, i)]
wait(futures)Uploading 8 chunks in parallel gives 8x upload throughput.
Download Flow
For download, reconstruct the file from chunks:
- Client requests
GET /files/{file_id}/download - Server returns the ordered chunk list (from
file_versions) - Client downloads chunks from CDN in parallel
- Client reassembles locally
CDN caches popular chunks at edge nodes. S3 is the origin store.
Sync Across Devices
When you edit a file on your laptop, your phone and tablet should see the update. Sync is the hard part.
Change log (event sourcing): Every file system operation is appended to a per-user change log:
user_changes table:
| change_id | user_id | timestamp | operation | file_id | metadata |
Operations: CREATE, UPDATE, DELETE, MOVE, SHARE
Sync protocol:
- Client sends its
last_sync_change_idto the server - Server returns all changes since that ID for this user
- Client applies changes locally (with conflict detection)
- Client sends its own local changes to the server
Notification:
Clients maintain a WebSocket or SSE connection. When a change is committed on the server, push a SYNC notification to all connected devices for that user. The device then fetches the diff.
User uploads file on laptop
↓
Server commits change, appends to change log
↓
Notification Service fetches all connected devices for this user_id
↓
Push "sync" notification via WebSocket
↓
Phone and tablet reconnect and fetch changes
Conflict Resolution
Two devices edit the same file while offline. Both push their versions when they reconnect.
Strategy 1 — Last Write Wins: Whoever commits later wins. Simple but lossy.
Strategy 2 — Duplicate and notify: Both versions are kept. User sees "Conflicting copy" (Dropbox does this). No data loss.
Strategy 3 — Operational Transform (Google Docs style): Transform concurrent operations so both can be applied. Works for text but complex to implement.
For files (not collaborative documents), Strategy 2 is the right choice. For text documents, OT or CRDT is needed.
Storage Architecture
With 10 exabytes of storage, no single file system works. Use distributed blob storage:
- S3 / GCS for raw chunk storage
- Organized by
chunk_hashprefix (first two chars) to distribute across buckets - Replicate across 3+ availability zones for durability
- Lifecycle policies: move infrequently accessed files to cheaper tiers (S3 Glacier)
s3://drive-chunks/ab/abcdef123456.../ ← chunk by hash prefix
s3://drive-chunks/cd/cdef789012.../
Erasure coding (used by Google Colossus): instead of 3 full replicas, divide a block into N shards and K parity shards. Can reconstruct from any N shards. 1.5x storage overhead vs 3x for full replication.
Access Control
Check permission before every read/write:
read_file(user_id, file_id):
owner = get_owner(file_id)
if owner == user_id: return file
share = get_share(file_id, user_id)
if share and share.permission in ['viewer', 'editor']:
if share.expires_at > now():
return file
raise PermissionDenied
Cache permission checks in Redis — they're read-heavy and the data is small.
For shared folders, permissions propagate down the file tree. Use a bitmask or role per path prefix for efficient checking.
Revision History
Keep the last N versions (e.g., 100) or versions from the last 30 days:
file_versions(file_id, version_num, chunk_list, created_at)
Since chunks are deduplicated, keeping 100 versions doesn't mean 100x storage. If only 3 chunks changed, only 3 new chunks are stored; all others are shared references.
Version cleanup: a background job deletes old file_versions rows and decrements ref_count on chunks. Chunks with ref_count = 0 are garbage collected.
Scaling Numbers
| Component | Scale | | -------------- | ----------------------------------------------------------------- | | Metadata DB | 1B files, sharded by user_id across 20 MySQL shards | | Blob Store | 10 exabytes → S3 handles this natively | | CDN | Cloudflare / Fastly for download throughput | | Upload Service | 100M DAU × avg 1 upload/day × avg 1MB = 100TB/day ingestion | | Sync events | 100M DAU × 5 changes/day = 500M events/day → Kafka handles easily |
Key Takeaways
- Chunk files for reliable uploads, parallel transfers, and deduplication
- Content-addressable storage (SHA256 hash as key) enables block-level deduplication
- Metadata and blob storage are separate — metadata in relational DB, blobs in S3
- Sync is event-sourced — a change log per user drives cross-device synchronization
- WebSocket / SSE for real-time sync notifications; fall back to polling
- Conflict resolution: keep both copies (simpler) or CRDT (collaborative text)
- Version history is cheap with deduplication — only changed chunks use new storage