TT Lab
Get started
Learn Learning paths Courses

Apache Hadoop — Stand up and run HDFS and YARN in one pod

Small files eat the NameNode's memory, not the disk

Continue in TT Lab

In one line

The cost of one file in HDFS is priced not in bytes but in the number of objects in NameNode memory. 2,000 files of 1KB are 2MB in terms of disk, but to the NameNode they are the objects of 2,000 files and 2,000 blocks. A Hadoop Archive (HAR) folds those objects into a few and makes the namespace lighter, and deleting the originals is up to a person.

Why this becomes a problem

The HDFS design document says the NameNode keeps the entire namespace and the block map in memory. Disk grows as much as you like by adding DataNodes, but the namespace lives inside the heap of a single NameNode process. One file, one directory and one block are all objects in that heap. So HDFS capacity has two axes: disk capacity measured in bytes, and namespace capacity measured in number of objects.

The same document makes clear that HDFS is tuned for large files. A typical file is gigabytes to terabytes, and it says one instance should support tens of millions of files. The scale of tens of millions is the size the design has in mind. Even with the same tens of millions of files, if the files average 1GB they hold tens of petabytes, but if they average 10KB, the namespace fills first while holding a few hundred gigabytes. This is also why the Federation documentation directly cites "deployments with many small files" as a reason to have several NameNodes.

Small files are a loss on the computation side too. According to the MapReduce tutorial, the number of maps is governed by the number of input blocks, and since preparing a task takes time, it is good for a map to run for at least a minute. 2,000 files of 1KB are 2,000 blocks, and map tasks that take a few seconds each to prepare start thousands of times to do work that takes less than a second.

How it works

The number of objects can be seen as numbers. In the FSNamesystem entry of the metrics documentation, FilesTotal is the number of files and directories and BlocksTotal is the number of allocated blocks. They can be read directly at /jmx of the NameNode web address (default port 9870). If you upload 2,000 small files into one empty directory, the two values move like this (the numbers in front are examples).

put 전     FilesTotal=  41   BlocksTotal=  12
put 2,000  FilesTotal=2042   BlocksTotal=2012   <- 파일 2,000 + 디렉터리 1, 블록 2,000

One small file creates one block. The default of dfs.blocksize in hdfs-default.xml is 128MB, but that is only the block's maximum size. The block of a 1KB file takes only a little over 1KB on the DataNode disk. Disk is not wasted. What is wasted is the NameNode's object that remembers that one block, and the work of reporting and managing that block. If you explain the small files problem as "a waste of disk", it becomes a wrong explanation.

A Hadoop Archive (HAR) folds these objects. According to the Hadoop Archives guide, a .har is one directory on the file system, and inside it are the metadata _index and _masterindex and the data files part-*. The contents of the original files are concatenated into a few part files, and the _index remembers the name of each original file and its position inside the part files. 2,000 files become a few files to the NameNode.

hadoop archive -archiveName logs.har -p /data/raw day1 /data/archive
hdfs dfs -ls har:///data/archive/logs.har/day1

A HAR exposes itself as a file system layer. Shell commands such as ls and cat work as they are with a har:// address, and it can also be used as MapReduce input. In exchange, there are several costs.

What a HAR reduces is the number of namespace objects. Even if you feed a HAR as MapReduce input, the logical files inside still look like files, so the problem of the number of maps is not solved by itself. If computation is the problem, it is better to merge into large files from the start.

Move and unpack with distcp

Bulk copying is done with DistCp. This is a MapReduce job too. It expands the list of files to copy and hands them out to the map tasks, and each map copies its own share. According to the documentation, it divides so that each map takes a similar number of bytes, but the file is the smallest unit, so if there are many small files, the copy is also dragged down by the number of files. -update copies only files that are absent or different at the target, so it is faster from the second run. Unpacking a HAR is also copying. If you hdfs dfs -cp with a har:// address as the source, it is unpacked sequentially, and with distcp, in parallel.

What it looks like in the field

First, the NameNode stalls on GC. Half the disk is empty, but the NameNode heap is full and long GCs repeat. Tracing the cause, it is a collector that drops thousands of small files every 5 minutes. This is why the trend of FilesTotal should be put in the monitoring items.

Second, too many files in one directory. The default of dfs.namenode.fs-limits.max-directory-items is 1,048,576. A collector that keeps dropping into one directory will eventually hit this wall and writes fail.

Third, you made the archive but the objects did not decrease. You did not delete the originals. Check the archive by reading it with a har address, and then delete the originals. If you reverse the order, there is no way back.

Fourth, the real solution is on the writing side. A HAR is a tool for clearing what has already piled up. For what piles up anew, you have to bundle it at the collection stage or have a regular job that merges it into large files, or the problem grows again.

What really matters in practice

What you will do in the next lab

You upload 2,000 small files to HDFS and measure the number of objects the NameNode has taken on by counting the files, directories and blocks with count and fsck. Without turning on YARN, you run hadoop archive in local mode to bundle those files into one HAR of replication 1, look at the listing with a har address to count the files inside, and then read one file and pull it out locally. You compare the object counts of the original directory and the HAR side by side and confirm in numbers how much they decrease, and finally make a backup copy of the archive with distcp.