System Design

Design Google Drive

September 10, 2024

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):

  1. Client initiates upload: POST /upload/initiate → returns upload_id
  2. Client queries which chunks are missing: GET /upload/{upload_id}/status
  3. Client uploads missing chunks in parallel: PUT /upload/{upload_id}/chunk/{index}
  4. 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:

  1. Client requests GET /files/{file_id}/download
  2. Server returns the ordered chunk list (from file_versions)
  3. Client downloads chunks from CDN in parallel
  4. 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:

  1. Client sends its last_sync_change_id to the server
  2. Server returns all changes since that ID for this user
  3. Client applies changes locally (with conflict detection)
  4. 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_hash prefix (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
VA
Vishal
Aggarwal

Full Stack Developer

Ask about Vishal ✦