Big Data Technologies

HDFS Architecture, Blocks, NameNodes, DataNodes and Commands

PGCP-BDA

HDFS

The Hadoop Distributed File System stores large files as replicated blocks across DataNodes while a NameNode maintains the filesystem namespace and block.

NameNode

The HDFS NameNode owns namespace metadata, permissions and the mapping from files to blocks

DataNode

An HDFS DataNode stores block replicas on local volumes, serves client reads and writes.

HDFS block

An HDFS block is a large fixed-size allocation unit independently stored and replicated, allowing different portions of a file to reside on different nodes.

replication factor

The replication factor is the desired number of HDFS block copies

filesystem namespace

The HDFS namespace is the hierarchy of directories, files, permissions, quotas and file-to-block metadata managed by the NameNode.

fsimage and edits

fsimage is a checkpoint of HDFS namespace state, while the edits log records later metadata transactions that are replayed during NameNode startup.

safe mode

HDFS safe mode is a NameNode state that permits metadata inspection but blocks ordinary namespace changes until enough block replicas have been reported.

HDFS command

An HDFS shell command uses the FileSystem client to request namespace operations from the NameNode and stream file data to or from DataNodes.

write once read many

HDFS is optimized for files written sequentially and read repeatedly

Distributed Storage

Why Distributed Storage?

  • Single machine cannot handle PB-scale data
  • Traditional storage: 1 machine; limited capacity
  • Distributed storage: Spread data across hundreds or thousands of machines
Single Server:          Distributed Storage:
+------+               +----+ +----+ +----+ +----+
|      |               |Node| |Node| |Node| |Node|
| 10TB |               | 1  | | 2  | | 3  | | 4  |
|      |               | 4TB| | 4TB| | 4TB| | 4TB|
+------+               +----+ +----+ +----+ +----+
                         Total: 16TB (easily expandable)

HDFS — Hadoop Distributed File System

HDFS is the distributed storage system designed for Big Data (part of Hadoop ecosystem).

Key concepts:

ConceptDescription
NameNodeMaster server; stores metadata (file locations, block info)
DataNodeWorker servers; store actual data blocks
Block sizeDefault 128MB per block (old: 64MB)
ReplicationDefault 3 copies of each block (fault tolerance)
File: sales.csv (500MB)

HDFS splits it into blocks:
Block 1 (0-128MB):   Stored on DataNode 1, 4, 7 (replication factor 3)
Block 2 (128-256MB): Stored on DataNode 2, 5, 8
Block 3 (256-384MB): Stored on DataNode 3, 6, 9
Block 4 (384-500MB): Stored on DataNode 1, 5, 7

Cloud Storage Services

ServiceProviderDescription
Amazon S3AWSObject storage; infinite scale; 99.999999999% durability
Azure Data Lake StorageAzureEnterprise analytics storage
Google Cloud StorageGCPObject storage with analytics integration
Google Cloud BigtableGCPWide-column NoSQL; millions of ops/sec

Measuring Performance

Measure representative workloads with realistic data volume, distribution, concurrency and parameter values. Record latency percentiles, rows examined, I/O, lock waits, CPU and execution frequency.

A query executed once per day can tolerate a plan that would be unacceptable thousands of times per second. A microbenchmark on an empty development table does not predict production behavior.

Change one design assumption at a time, re-check plans and include write cost in the result.

EXPLAIN

EXPLAIN shows the chosen plan:

EXPLAIN SELECT employee_id, salary FROM employee WHERE department_id = 10 AND salary >= 50000;

Important fields can include access type, possible keys, selected key, key length, estimated rows, join order and extra operations.

Read the plan as a whole. Seeing an index name does not prove efficiency; it might scan the entire index, examine many rows, sort a large intermediate result or perform costly lookups.

Pagination

Offset pagination:

ORDER BY ordered_at DESC, order_id DESC LIMIT 20 OFFSET 100000

may scan and discard many prior rows.

Keyset pagination continues after the last seen key:

WHERE (ordered_at, order_id) < (?, ?) ORDER BY ordered_at DESC, order_id DESC LIMIT 20

With a matching index, cost stays closer to page size. It also avoids shifts caused by newly inserted earlier rows. The ordering must include a unique tie-breaker.

Why a Table Scan Can Be Correct

An index lookup involves traversing the index and possibly performing many primary-row lookups. When a query returns a large percentage of the table, sequentially scanning pages can be cheaper.

Small tables also make scans inexpensive. A scan is not automatically a performance defect.

Judge the plan from rows examined, I/O pattern, latency, concurrency and workload frequency. Forcing an index without evidence can make a query slower and more fragile as data distribution changes.

Sorting and Grouping

An index can supply rows in an order compatible with ORDER BY and avoid a separate sort:

WHERE customer_id = ? ORDER BY ordered_at DESC, order_id DESC

A matching index might begin with customer_id, ordered_at, order_id in compatible directions and version-supported form.

If filtering and ordering requirements conflict, the optimizer chooses a tradeoff. LIMIT makes ordered index access especially valuable when it can stop early.

Only ORDER BY defines the result-order contract, even if the selected plan happens to scan an index.

Continue learning

Related notes

Put this topic into timed practice

Open mock tests when you want full-exam pacing, or keep drilling in practice mode.