NEW: ML Mock & Coaching now available

Questions

Distributed File System

MetaGoogleDatabricks

Design a S3-like distributed file system - This system design covers hierarchical namespace, block storage, strong consistency, and petabyte-scale durability from the ground up.

45 min read

Challenge

Think Beyond the Happy Path

Distributed file systems aren't just about storing bytes — they're about petabyte-scale durability, strong consistency for metadata, and reducing read latency across globally distributed replicas.

Before diving into the material, take a moment to ask yourself:

  • Do you know how to reduce read latency at scale — what caching strategies and replica placement would you use for a globally distributed file system?
  • Do you know how to ensure strong consistency for file and directory operations — what happens when two clients try to create the same file simultaneously?
  • Do you know how to guarantee no data loss — how many replicas do you need, and how do you handle the failure of multiple storage nodes?
  • Do you know how the system supports high scalability — how do you shard metadata and avoid hotspots when millions of files are in the same directory?

You don't need to answer all of these right away. A Senior/Staff+ Engineer doesn't stop at basic functionality — they anticipate edge cases, design for resilience, and push for production-grade reliability.

But great systems start with great questions. What would you ask next?

Problem Statement

Design a S3-like distributed file system without needs on AWS features from the ground.

What is AWS S3?

Amazon S3 (Simple Storage Service) is a scalable, distributed object storage system widely used for storing and retrieving large amounts of data, including files, backups, media, and logs. It organizes data into buckets and supports features like versioning, lifecycle policies, and multi-region durability. S3 is optimized for durability, availability, and throughput at internet scale, and it abstracts away complex infrastructure, making it a common mental model when designing file storage systems in interviews.

Functional Requirements

FR1 – createDirectory(path)

Create new directory.

FR2 – deleteDirectory(path)

Delete a file or empty directory.

FR3 – createFile(path, data)

Create a new immutable file.

FR4 – readFile(path, offset, length)

Retrieve (partial) contents of a file.

FR5 – listDirectory(path)

List children of a directory.

Non-Functional Requirements

NFR1 – Scalability

100s of PBs & millions of directories.

NFR2 – Consistency

Strong for directory metadata ops, eventual for block replicas.

NFR3 – Durability

No data loss; redundant in block storage.

NFR4 – Latency

Low latency in file read and directory listings.

Below the line:

  • Multi-tenancy - Customers logically isolated
  • Extensibility - Potential for versioning, ACLs, quotas later
Constraints & Assumptions

To frame the design properly, let's break down the problem space through the key constraints & assumptions given in the interview prompt. Each of them not only narrows down our choices but also tells us what matters most in this system.

  1. Immutable Files

    Once a file is written, we assume it can't be modified—only read or deleted. This assumption simplifies consistency and concurrency handling, but it raises the bar for durability. Every file write must ensure the data is fully persisted, verified, and recoverable, as there's no "fix it later" path.

  2. Mutable Directories

    Users can add or delete files and subdirectories at any time, so the directory structure must support concurrent mutations with strong correctness guarantees. This is where transactional integrity and concurrency control (e.g., locks or database semantics) come into play.

  3. Use Only Primitive Infrastructure Building Blocks

    We're not allowed to use ready-made distributed storage like S3 or GFS. Instead, we can only build using low-level components like:

    • RDBMS (e.g., PostgreSQL) for relational, transactional metadata

    • KV Stores (e.g., Redis or RocksDB) for fast, hashed lookups

    • Zookeeper for distributed coordination (locking, routing)

    • Single-machine filesystems as the physical block store

      This constraint tests your ability to architect a reliable distributed system from first principles.

Tip

Apache Zookeeper is a lightweight coordination service designed for distributed systems. It works like a consistent, replicated key-value store with a file-like hierarchy, where each node (called a znode) can store data and metadata. It guarantees strong consistency using a leader-based protocol (ZAB), making it ideal for managing distributed locks, leader election, and synchronized access. In our file system design, Zookeeper helps safely coordinate directory-level operations across multiple nodes, ensuring correctness when many users modify paths concurrently, which we will cover later.

Scale & Region Assumptions
  1. Scale: Millions of Users, 100s of Petabytes

    We're targeting massive scale: billions of files and directories spread across petabytes of data. This mandates horizontally scalable architecture, efficient metadata sharding, and smart data placement strategies to avoid bottlenecks.

  2. Single Region & Data Center

    The system only needs to run in a single region and data center. This removes cross-region latency and replication complexity from our current scope, letting us focus more on core file system design. (But note: multi-region could be an extension question.)

Common Pitfalls to Avoid

Before diving further into the design, many candidates make critical mistakes such as:

  • Overusing abstract cloud terms like "just like S3" without showing how it's built from scratch.
  • Ignoring immutability constraints, especially when handling file writes or updates.
  • Conflating file and directory semantics, leading to incorrect handling of metadata operations.
  • Underestimating metadata bottlenecks — the system's scalability often hinges more on managing trillions of paths than raw data volume.
  • Skipping consistency guarantees — assuming eventual consistency is "good enough" everywhere, even for hierarchical operations.

Requirement Summary

Functional Requirements (FRs)
NameDescription
1. createDirectory(path)Create new directory.
2. deleteDirectory(path)Delete a file or empty directory.
3. createFile(path, data)Create a new immutable file.
4. readFile(path, offset, length)Retrieve (partial) contents of a file.
5. listDirectory(path)List children of a directory.
Non-Functional Requirements (NFRs)
NameDescription
1. Scalability100s of PBs & millions of directories.
2. ConsistencyStrong for directory metadata ops, eventual for block replicas.
3. DurabilityNo data loss; redundant in block storage.
4. LatencyLow latency in file read and directory listings.

Core Entity

  • file (immutable, content-hashed, block-indexed)
  • directory (mutable, hierarchical, maps path -> children)
  • block (fixed size chunk(64MB), deduplicated by hash)
  • path (absolute string (/abc/def/file))
  • customer (tenant root-level entity(namespace, isolation))

APIs - with Entity Mapping

  • createDirectory(path: str) -> boolean

    Inserts a new directory into namespace; updates path hierarchy under customer.

  • deleteDirectory(path: str) -> boolean

    Deletes a directory from namespace after verifying no file or sub-directory children exist.

  • createFile(path: str, data: bytes) -> boolean

    Creates a file entry, stores content as hashed blocks, links metadata under directory.

  • readFile(path: str, offset: int, length: int) -> bytes

    Reads file metadata and fetches corresponding block ranges by block_hash.

  • listDirectory(path: str) -> list[str]

    Lists all file and directory entries under given path in namespace.

High Level Design

High Level Design Diagram

Click to expand

Building Blocks

To handle hundreds of petabytes of storage and trillions of files efficiently, we break the system into modular components, each with a specific role in the file system pipeline. This separation of responsibilities improves scalability, clarity, and fault isolation. Below is a high-level overview of the core components and how they fit together.

API Gateway

The API Gateway is the entry point for all external requests. It exposes high-level file system APIs like putFile, getFile, createDirectory, and others. It is stateless and horizontally scalable, simply forwarding validated requests to the backend service layer. This layer may also handle basic normalization, authentication, and retry logic.

File Service

The File Service is the central orchestrator for file system operations. It coordinates with other components to handle file uploads, reads, and metadata changes. For example, on a file upload, it splits the file into chunks, stores them across storage nodes, updates metadata, and commits the full file atomically. It ensures the correctness and consistency of user-visible operations.

Metadata Store

The Metadata Store tracks the structure of the file system—what paths exist, which are directories or files, who owns them, and which blocks make up a file. It stores two main types of metadata: directory hierarchy and file-to-block mappings. This store is strongly consistent and supports transactional updates, making it the foundation of the system's correctness.

Chunk Manager

The Chunk Manager maintains a fast-access mapping from block hash to the list of nodes that store that block. It acts like a distributed address book for locating content-addressed data. It is backed by a key-value store and optimized for high-throughput lookups during file reads and uploads. It may also track block usage for garbage collection.

Block Store

The Block Store defines the logic for handling immutable blocks of data. Files are broken into fixed-size blocks (e.g., 64MB), each addressed by a SHA256 hash. These blocks are never changed once written. The block store handles integrity checks and provides APIs to read or write a specific block by hash.

Storage Nodes

Storage Nodes physically store and serve the blocks on local disk. Each node handles basic upload and download operations for the blocks it holds. Nodes are simple by design and scale out horizontally. They report their status to the chunk manager and are periodically cleaned up by background workers.

Garbage Collection Workers

Garbage Collection (GC) Workers clean up unused data and metadata asynchronously. When a file is deleted or no longer referenced, GC workers remove its associated blocks from storage and update metadata entries accordingly. These workers run in the background and are carefully throttled to avoid disrupting live traffic.

Sign in to continue reading

"Distributed File System" requires a free account to access.

Sign in to continue