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.
-
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.
-
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.
-
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.
-
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
-
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.
-
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) | |
|---|---|
| Name | Description |
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) | |
|---|---|
| Name | Description |
1. Scalability | 100s of PBs & millions of directories. |
2. Consistency | Strong for directory metadata ops, eventual for block replicas. |
3. Durability | No data loss; redundant in block storage. |
4. Latency | Low 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) -> booleanInserts a new
directoryinto namespace; updatespathhierarchy undercustomer. -
deleteDirectory(path: str) -> booleanDeletes a
directoryfrom namespace after verifying nofileor sub-directorychildren exist. -
createFile(path: str, data: bytes) -> booleanCreates a
fileentry, stores content as hashedblocks, links metadata underdirectory. -
readFile(path: str, offset: int, length: int) -> bytesReads
filemetadata and fetches correspondingblockranges byblock_hash. -
listDirectory(path: str) -> list[str]Lists all
fileanddirectoryentries under givenpathin namespace.
High Level Design

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