Apache Hadoop — Stand up and run HDFS and YARN in one pod
Small files eat the NameNode's memory, not the disk
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.
- It is immutable. As the documentation says, renaming, deleting and creating inside a HAR produce an error. To fix one day's worth, you recreate the archive.
- A MapReduce job runs to create it. It usually runs on YARN, and if YARN is not turned on, you run it inside the client JVM with
-D mapreduce.framework.name=local. distcp is the same. - It does not delete the originals. As the documentation stresses, to reduce the namespace you have to delete the originals yourself after checking the archive. Until you delete them, the number of objects has actually increased.
- The default replication is 3. If you do not give
-r, it creates with replication 3. In this lab environment with one DataNode, if you do not give-r 1, under-replicated blocks appear. - Reading goes through one more step. To read one file, it finds the position in the index and then reads that range of the part file. It is for archival, not for reading a little at a time frequently and at random.
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
- The cost of a file is the number of objects. Files, directories and blocks are all objects in the NameNode heap.
- Small blocks do not waste disk. What is wasted is the NameNode's memory and the number of tasks.
- Look at FilesTotal and BlocksTotal. If you look only at disk utilization, this problem is invisible.
- A HAR does not delete the originals. Objects decrease only when you delete them yourself after checking.
- A HAR is immutable and its default replication is 3. It does not suit data that needs fixing, and on a small cluster you give -r.
- The fundamental solution is to bundle on the writing side. You cannot plug the leak with a tool that clears up.
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.