Distributed File System (HDFS / GFS / Ceph)
PremiumDistributed File System case study covering HDFS, GFS, and Ceph patterns: master/metadata server (NameNode HA via QJM + ZooKeeper) plus DataNode tier with rack-aware 3x replication. POSIX-like file API contrasted with object storage. Block-based storage (128MB), pipelined writes, streaming reads, append-only semantics. Five scenarios: pipelined chunk write to 3 replicas, parallel multi-block read with Master out of data path, append to existing file, DataNode failure with background re-replication, Active NameNode crash with Standby promotion via fence tokens. Includes 2 ADRs (single master vs sharded vs masterless Ceph CRUSH; 3x replication vs Reed-Solomon EC) and capacity hints for all nodes.
Что внутри
Distributed File System: HDFS HA Reference
Scope
This diagram is specifically an Apache HDFS high-availability reference design using the Quorum Journal Manager (QJM). It does not combine incompatible claims from HDFS, Google File System (GFS), and Ceph.
HDFS uses a NameNode for the filesystem namespace and block placement, while file bytes travel directly between clients and DataNodes. GFS has a primary/secondary mutation protocol and record-append semantics described in its paper; those are not labels for a normal HDFS pipeline. Ceph's CRUSH helps place RADOS objects without a central allocation table, but CephFS still uses Metadata Servers for its filesystem namespace. “All distributed filesystems are masterless” is therefore false.
Metadata invariants
- FsImage/EditLog persist the namespace, including file-to-block relationships and replication policy. Current DataNode block locations are reconstructed and refreshed from heartbeats/block reports; they are not immutable location rows in FsImage.
- Exactly one NameNode is Active. The Active durably logs each namespace edit to a majority of JournalNodes. With five JNs, the majority is three and normal progress tolerates at most two unavailable JNs.
- The Standby tails committed edits and receives block reports/heartbeats from DataNodes. Before promotion it reads all committed journal edits.
- JournalNodes allow one writer, protecting the edit log from split-brain writers. Configured fencing is still required to stop a previous Active from serving stale reads or causing external side effects.
- Clients use a logical HA nameservice/failover proxy and retry operations according to operation semantics. A timeout alone is not proof that a mutation did not commit.
QJM is described here as a majority replicated edit journal with writer exclusivity, not hand-waved as “Paxos-like.” ZooKeeper/ZKFC coordinates automatic failover; the diagram does not claim that every data block is written through ZooKeeper or the journals.
Полный разбор, ADR-ы, сценарии и deep dives — после оплаты бандла.
System Design Cases
Полный доступ ко всем кейсам бандла
Premium открывает полный разбор для подготовки к интервью
- Где архитектура ломается первой и как защищать выбранный дизайн.
- Конкретный capacity math: размеры данных, throughput и пороги масштабирования.
- Trade-off-ы в стиле ADR, которые легко превращаются в структурированный ответ.
- Запускаемые сценарии: happy path, отказы, retry и recovery.
Регистрация бесплатна. Оплата — следующим шагом, из этого же кейса.
Уже есть аккаунт?