Skip to content

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

8 Commits

Folders and files

Repository files navigation

Google File System (GFS) MVP Replica

This repository contains a solo-built, distributed storage prototype inspired by the Google File System architecture. It is implemented in Python using gRPC for RPC communication, Docker Compose for local deployment, and a simplified metadata/data separation model.

Project Summary

The system is organized as a single-master, multi-chunkserver distributed file system:

  • Master Node manages metadata, file namespaces, chunk allocation, leases, and liveness monitoring.
  • Chunkservers store raw chunk data on local volumes and execute read/write/append requests.
  • Client Library routes file operations through the Master and directly to chunkservers for data plane access.

This project demonstrates distributed systems design patterns such as leader metadata coordination, replication, failure detection, lease-based primary assignment, and write forwarding.

File Structure

.
├── DESIGN.md
├── README.md
├── docker-compose.yml
├── docker/
│   ├── Dockerfile.chunkserver
│   └── Dockerfile.master
├── master/
│   ├── master_node.py
│   └── namespace_manager.py
├── chunkserver/
│   └── chunk_server.py
├── client/
│   └── client_lib.py
├── proto/
│   ├── gfs.proto
│   ├── gfs_pb2.py
│   ├── gfs_pb2_grpc.py
│   └── gfs_pb2.pyi
├── gfs_pb2.py
├── gfs_pb2_grpc.py
└── benchmarks/

Key components

  • master/master_node.py: Master metadata service implementation.
  • master/namespace_manager.py: In-memory hierarchical namespace manager with directory auto-creation.
  • chunkserver/chunk_server.py: Chunkserver storage engine and replication coordinator.
  • client/client_lib.py: Example client library for create/read/write/append operations.
  • proto/gfs.proto: gRPC service definitions and message schema.
  • docker-compose.yml: Local cluster orchestration with one master and five chunkservers.

Architecture & Design

Master Node

The master performs metadata-only duties:

  • Maintains a directory tree and file metadata with NamespaceManager.
  • Allocates new chunks for files and chooses replica chunkservers.
  • Grants primary leases for chunk writes and increments chunk versions.
  • Receives heartbeats from chunkservers and marks nodes DEAD if they miss 3 intervals.
  • Persists metadata operations to an append-only operation log at /data/master_op_log.jsonl.
  • Implements lazy garbage collection by renaming deleted file entries to a hidden .deleted/ namespace.

Chunkserver

Each chunkserver:

  • Stores fixed-size chunks of 1 MB in /data.
  • Handles ReadChunk, WriteChunk, AppendRecord, and ReplicateChunk RPCs.
  • Uses per-chunk locks to serialize access and avoid corrupting concurrent writes.
  • Registers itself automatically with the master at startup.
  • Sends periodic heartbeats with chunk version reports and free-space simulation.
  • Includes a boot scan that detects existing local chunks and initializes their version tracking.

Client Library

The client library demonstrates the control/data-plane split:

  • Creates file metadata through the master.
  • Retrieves chunk location information from the master.
  • Sends direct read/write requests to the primary chunkserver.
  • For writes, forwards mutation requests to secondaries for replication.
  • Maps container network addresses to host-local ports so clients can run outside Docker.

RPC Protocol

The project uses proto/gfs.proto to define two services:

  • MasterService: register chunkservers, heartbeat, file create/get chunk locations/delete.
  • ChunkService: read/write chunks, append records, and replicate chunk content.

How to Run

Prerequisites

  • Docker Engine
  • Docker Compose
  • Python 3.11 (for local development and generated stubs)

Start the Cluster

From the repository root:

docker compose up --build

This command builds both the master and chunkserver images and starts:

  • master on port 50051
  • chunkserver1 on port 50052
  • chunkserver2 on port 50053
  • chunkserver3 on port 50054
  • chunkserver4 on port 50055
  • chunkserver5 on port 50056

Inspect Running Services

docker compose ps

Simulate Failure

Stop a chunkserver to observe heartbeat and liveness detection:

docker compose stop chunkserver2

Restart it when ready:

docker compose start chunkserver2

Client Usage Example

The project does not include a full application, but you can interact with the system using client/client_lib.py.

Example Python usage:

from client.client_lib import GFSClient

client = GFSClient(master_address="localhost:50051")
client.write("/example/file.txt", b"Hello GFS", offset=0)
print(client.read("/example/file.txt", offset=0, length=32))

Note: The client library translates internal container addresses such as chunkserver1:50052 to host ports localhost:50052, so this example works against the Docker Compose cluster.

Design Decisions

  • Separation of metadata and data path: Master handles only namespace and chunk placement, while chunkservers handle actual read/write data.
  • Lease-based primary assignment: The master grants a primary lease for each chunk and tracks a lease expiration time to avoid stale writers.
  • Synchronous replication on write: The primary forwards each write to configured secondaries for consistency.
  • Heartbeat-based failure detection: Chunkservers send periodic heartbeats; the master evicts nodes that stop reporting.
  • Operation log persistence: Master metadata operations are appended to a durable JSONL log so the namespace can be replayed after restart.
  • Lazy deletion: Deleted files are retained in a .deleted/ namespace first, which simplifies safe removal and deferred cleanup.
  • Boot-time chunk scan: Chunkservers scan existing local storage on startup to recover persisted chunks and avoid silent data loss.

Problems Faced

  • Cluster addressing inside Docker: The system needed address translation because Docker container hostnames differ from host-local ports.
  • Concurrent chunk mutations: Per-chunk locks were required to prevent concurrent writes and appends from corrupting chunk files.
  • Stale replica detection: Master version tracking and heartbeat reporting had to be combined to detect out-of-date chunk replicas.
  • Recovery semantics: Replication and primary write coordination had to be implemented carefully so secondaries receive the same data.
  • Operation replay correctness: Master metadata persistence had to recover file creation and deletion state on restart.

Limitations & Scope

This implementation is a minimum viable prototype and intentionally omits features such as:

  • Multi-chunk files and chunk indexing beyond a single chunk per file
  • Real distributed consensus or leader election
  • Full chunk re-replication and automatic recovery after node failure
  • Access control and authentication
  • Complete directory listing and namespace query APIs

What This Project Demonstrates

  • Distributed system architecture design and service decomposition
  • gRPC protocol design and Python service implementation
  • Docker Compose orchestration for a multi-node prototype
  • Fault-detection and basic replication strategies
  • Resume-ready storytelling around tradeoffs, failure handling, and system behavior

Next Steps / Improvements

Potential improvements to highlight on your resume:

  • Add support for multi-chunk files and chunk striping.
  • Implement a real replica recovery pipeline with chunk re-replication.
  • Add stronger consistency guarantees with consensus and versioned leases.
  • Expose a full filesystem API with directory listing and rename semantics.
  • Add integration tests and end-to-end workload benchmarks.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages