I built a small distributed file system in Go, modeled loosely on HDFS. It has three binaries:
- a controller that owns metadata: which nodes are alive and which nodes hold which chunks
- storage nodes that keep chunk bytes on local disk
- a client with
store,get,getmeta,delete,ls, andnodescommands
Everything talks over plain TCP using Protocol Buffers. There is no gRPC and no external coordination service. This post walks through how the pieces fit together and the failure cases I had to handle.
The source code is on GitHub: Shreyas-Yadav/hdfs-lite.
The wire protocol: one envelope, length-prefixed
Every message type is a field in a single oneof called Wrapper, so each connection carries one kind of thing: a Wrapper. That covers RegisterNode, StoreChunkRequest, Heartbeat, ReplicateChunkRequest, CommitFileRequest, and the rest.
TCP is a byte stream, so each message needs a boundary. MessageHandler writes a 4-byte big-endian length followed by the protobuf body. On the read side, ReadMsg reads 4 bytes, then reads exactly that many bytes with io.ReadFull, then unmarshals. Callers just call Send(m), which type-switches on the message and wraps it. On the receive side they use typed helpers like ReceiveStoreChunkResponse(), built on a small generic receiveTyped[T], which return an error if a different message type shows up.
Each server accepts connections in a loop and starts one goroutine per connection. That goroutine loops over ReadMsg and dispatches on msg.Msg.(type).
Joining the cluster and staying alive
On startup, a storage node binds its listening port first and registers with the controller only after that, so it never announces an address it can't serve. The node ID is simply host:port. After registering, the node keeps the TCP connection to the controller open and sends a Heartbeat over it every 5 seconds.
Each heartbeat carries the node's free space (from syscall.Statfs), a count of requests handled, and a new_chunks list of chunks stored since the last beat. If a heartbeat fails, reconnectLoop tries to re-register. It waits 5 seconds first and doubles the wait after each failure. It stops doubling once the wait reaches the 60-second maxBackoff, so the longest wait in practice is 80 seconds.
On the controller, monitorNodes runs every monitorInterval (5s). Any node that hasn't sent a heartbeat within nodeFailureTimeout (15s, about three missed beats) is marked dead, and that triggers a replication repair pass. If a dead node later sends another heartbeat, the controller clears the dead flag and starts another repair pass.
Placing chunks
The client splits a file into fixed-size chunks and sends a StoreFileRequest. The replication factor defaults to 3. The controller rejects the request if the file already exists or if fewer live nodes are available than the replication factor requires.
selectNodes handles placement with an offset round-robin over the live nodes, sorted by ID. The node list for chunk i starts at index i, so the primary role rotates across nodes and the write load is spread. The first node in each placement is the primary and the rest are replicas. This scheme ignores free space. Free space only matters later, when the controller picks a node to repair onto.
Writing: pipeline replication
The client doesn't send each chunk three times. It sends each chunk once, to the primary, with the rest of the chain in Replicas and a SHA-256 checksum of the data. The primary:
- checks that the checksum matches the received bytes
- writes the chunk plus a
.checksumsidecar file to disk - forwards the chunk to
Replicas[0]and passes along onlyReplicas[1:]
Each hop waits for the next hop's ack before it acks upstream. So when the client gets ok from the primary, every node in the chain has already written the chunk. If a forward fails, the node deletes its own copy before reporting failure:
if len(req.Replicas) > 0 {
next := req.Replicas[0]
if err := forwardChunk(next.Address, req); err != nil {
if rbErr := node.deleteChunk(req.Filename, req.ChunkIndex); rbErr != nil {
log.Error("rollback failed for chunk %d of %s: %v", req.ChunkIndex, req.Filename, rbErr)
}
_ = handler.Send(&pb.StoreChunkResponse{
Ok: false,
Message: fmt.Sprintf("pipeline forward to %s failed: %v", next.Address, err),
})
continue
}
}
Committing: no half-written files in the namespace
When the controller receives a StoreFileRequest, it puts the placement plan in pendingFiles, not files. The client uploads every chunk and then sends CommitFileRequest on the same controller connection. Only then does commitFile move the record into files and write state to disk.
If the client dies partway through, its controller connection closes. A deferred block in the connection handler sees a pending filename that never got committed and calls abortPending. That removes the record and asks every planned node to delete its chunk. On storage nodes, deleting a chunk that doesn't exist counts as success, so cleaning up a partial upload is safe.
Only committed files are persisted, as JSON in controller_state.json. Node membership is not persisted. The controller rebuilds it as nodes register and send heartbeats.
Reading: parallel fetch, positional writes
getFile gets the placements from the controller, creates the output file, and truncates it to its full size up front. It then fetches chunks in parallel. A buffered channel caps the fetches at 8 goroutines at a time:
sem := make(chan struct{}, 8)
var wg sync.WaitGroup
var firstErr error
var errMu sync.Mutex
for _, placement := range placements {
wg.Add(1)
sem <- struct{}{}
go func(p *pb.ChunkPlacement) {
defer wg.Done()
defer func() { <-sem }()
data, err := fetchChunkWithFallback(filename, p)
Each goroutine writes its chunk with WriteAt at offset ChunkIndex * ChunkSize. Chunks can finish in any order, and the file is never assembled in memory. fetchChunkWithFallback tries each replica in placement order and checks the SHA-256 of what it receives. If the data was corrupted in transit, or any other error occurs, it moves on to the next replica.
Self-healing reads on the storage node
Checksums are also checked on disk. When a storage node serves a GetChunkRequest, it compares the chunk against its .checksum sidecar. If the chunk is missing or the checksum doesn't match, the node doesn't just return an error. It calls repairChunkFromReplicas, which fetches a good copy from another replica, stores that copy locally, and then serves it.
Replica A might fetch from B while B's copy is also bad, and B would then try to repair from A. GetChunkRequest.avoid_node_ids prevents that loop. Each node adds itself to the list before asking the next one:
avoidSet[s.nodeID] = true
nextAvoidNodeIDs := append(append([]string(nil), avoidNodeIDs...), s.nodeID)
for _, placement := range placements {
if placement.ChunkIndex != chunkIndex {
continue
}
for _, replica := range placement.Nodes {
if replica == nil || replica.Address == "" || avoidSet[replica.NodeId] || replica.Address == s.address {
continue
}
data, err := fetchReplicaChunk(filename, chunkIndex, replica.Address, nextAvoidNodeIDs)
if err != nil {
continue
}
As a result, a client read also repairs a corrupted local copy.
Re-replication after a node dies
reconcileUnderReplicated brings chunks back up to their target replica count. It works in three phases so the controller mutex is never held during network calls:
- Step A (under lock): Find every chunk with fewer live replicas than its target and queue the repair work: a healthy source, a replacement node, and the dead node to replace.
-
Step B (no lock): For each work item, send a
ReplicateChunkRequestto the source storage node. The source verifies its own checksum and pushes the chunk straight to the destination. Chunk bytes never pass through the controller. - Step C (under lock): Swap the replacement into the placement and save state once at the end.
pickReplacementNode picks the live node with the most free space that doesn't already hold the chunk. If no spare node exists, the chunk stays marked under-replicated and a later pass retries it.
Why the 3-minute grace period
This one came from a real failure mode. When the controller restarts, it reloads every file's placements, but it doesn't know about any nodes yet. If it ran a repair pass right away, every chunk would look under-replicated, and the controller would start copying data around the cluster for no reason.
So the first pass waits for startupGracePeriod:
go controller.monitorNodes()
go func() {
time.Sleep(startupGracePeriod)
controller.reconcileUnderReplicated("controller startup")
}()
startupGracePeriod is 3 minutes. The repair triggers on node registration and on node recovery are also skipped during this window. That gives storage nodes, which retry on their own backoff, time to reconnect before the controller makes any replication decisions.
Deletes
The client sends DeleteFileRequest to the controller. The controller asks every replica to delete its chunk, then drops the metadata. If some replicas can't be reached, the delete still succeeds from the client's point of view, and the response message says how many replica deletes couldn't be confirmed.
What I'd change next
Some limitations I know about:
-
The controller is a single point of failure. Metadata goes to one local JSON file, and
saveStatewrites it withos.WriteFiledirectly, not a write-to-temp-then-rename. A crash during the write could corrupt the state file. -
Uploads are sequential per chunk. Reads run 8 at a time, but
storeFilereads and pushes chunks one after another. -
new_chunksis reported but barely used. Storage nodes send it in every heartbeat, but the controller only logs it. The controller never compares its placement map against what nodes actually have on disk. - Initial placement ignores capacity. Round-robin spreads load evenly, but a nearly full node gets as many new chunks as an empty one.
-
The liveness window is duplicated. Two functions hard-code
15*time.Secondinstead of usingnodeFailureTimeout.
The ideas I learned the most from were pipeline replication with rollback, commit-after-upload, and checksum-verified reads that repair the bad copy. None of them took much code, but each one closed a failure case that I only understood after it had happened.
Top comments (0)