Apache Hadoop — Stand up and run HDFS and YARN in one pod
A file is split into blocks and laid on DataNode disks
In one line
An HDFS file is cut into blocks of the same size and placed on DataNode disks as ordinary files, and each block gets as many copies as the replication factor. The block size can be set per file, but once set, it changes the number of blocks, the load on the NameNode, and even the value of the file checksum.
Why blocks are this large
A local file system's blocks are usually a few KB. An HDFS block is, per dfs.blocksize in hdfs-default.xml, 134217728 bytes by default, that is 128MB. There are reasons for a difference of tens of thousands of times.
First, HDFS was designed for reading large files from start to end. The design document says it values throughput over latency. With large blocks, you position once and then read sequentially for a long time. Second, each block is an object in NameNode memory. If you cut a 1TB file into 4KB blocks, there are well over 200 million blocks, and the NameNode would have to hold the locations of all of them in memory. Third, a block is also used as a unit for dividing computation. As the design document puts it, moving computation is cheaper than moving data, so engines such as MapReduce run the work where the block is.
That does not mean you can make it arbitrarily small. dfs.namenode.fs-limits.min-block-size is 1048576 bytes (1MB) by default, and its description says it is a lower limit to prevent a surge of blocks from accidentally setting a very small block size.
How it works
The block size differs per file. The design document says the block size and the replication factor can be set per file. The cluster default is used only for files for which you did not specify one at creation. Also, all blocks except the last one are the same size. So if you upload a 12MiB file with a 4MiB block, it becomes 3 blocks, and with the default 128MB, it becomes 1 block. Even a 1KB file takes one block. 50 small files are 50 blocks: blocks are not shared between files.
# 이 파일만 블록 4MiB 로 올린다
hdfs dfs -D dfs.blocksize=4194304 -put big.dat /data/big.dat
# 블록 목록과 위치를 본다
hdfs fsck /data/big.dat -files -blocks -locations
The fsck in the HDFS commands documentation shows, with -files -blocks -locations, the block IDs of each file and the DataNodes that hold those blocks. Once you have a block ID here, you can go down to the next level.
Blocks on disk. The place where a DataNode puts blocks is dfs.datanode.data.dir, which is file://${hadoop.tmp.dir}/dfs/data by default (the lab image has moved it to /var/lib/hadoop/data). If you follow it down, there are ordinary files starting with blk_, and next to them a paired file with the same name plus .meta. The former is the bytes of the block as they are, and the latter is the checksum of those bytes. The design document explains that a DataNode does not pile files into one directory but splits them into subdirectories on its own, because a local file system cannot handle a huge number of files in one directory well. If you uploaded a text file, you can open the blk_ file with head and read the original contents as they are. This means HDFS does not use a special storage format.
Replication. The default replication factor is dfs.replication 3. The design document explains the placement policy when the replication factor is 3 as follows: one on the writer's machine (or a random node on the same rack), one on a node on a different rack, and the last on another node on that other rack. It survives the loss of a whole rack, and reduces write traffic between racks.
There is one important constraint. The NameNode does not put more than one replica of the same block on one DataNode, so the maximum number of replicas that can be made is the number of DataNodes at that time. In a lab environment with only one DataNode, if you run setrep 3 from the file system shell, the command succeeds but the copy stays one forever. fsck reports such a block as under-replicated. A command succeeding does not mean replication happened.
Why the checksum is tied to the block size
HDFS verifies bytes in small pieces. dfs.bytes-per-checksum is 512 and dfs.checksum.type is CRC32C by default. One CRC is computed for every 512 bytes and kept in the .meta file.
The problem is the file-level checksum that hadoop fs -checksum returns. The default combine method dfs.checksum.combine.mode is MD5MD5CRC. As the name says, it binds the per-512-byte CRCs with an MD5 for each block, and binds those block MD5s again with an MD5. The block boundaries are part of the computation, so even with the same content, the value differs if the block size differs. The description in hdfs-default.xml also says that the original method cannot compare files with different block layouts, while methods such as COMPOSITE_CRC can compare regardless of block layout.
hadoop fs -checksum /a/4m.dat /a/128m.dat # 값이 다르다
hadoop fs -D dfs.checksum.combine.mode=COMPOSITE_CRC \
-checksum /a/4m.dat /a/128m.dat # 값이 같다
Where this becomes an accident in the field is copy verification. The DistCp documentation says -update compares the size, block size and checksum of the source and target. A copy moved without preserving the block size may appear to have a mismatched checksum under the default method.
What it looks like in the field
First, "I changed it to replication 3, but fsck keeps warning." It happens when there are fewer than three DataNodes, or no place that satisfies the rack placement. setrep only changes the target, and actually making the copies is the back-end work where the NameNode instructs the DataNodes.
Second, if you shrink the block size, the NameNode hurts first. The number of blocks is the file size divided by the block size, so halving the block size doubles the block objects. It is a burden in the same direction as having many small files.
Third, a different checksum on a copied file does not mean the data is corrupted. It may be because the block sizes differ. You can tell by changing the combine mode to COMPOSITE_CRC and comparing again.
Fourth, do not touch the blk_ files on a DataNode disk by hand. Those files are paired with the NameNode's block map. If you move or delete them, the NameNode learns late through a block report, and in between, reads fail.
What really matters in practice
- Block size and replication factor are per-file attributes. The cluster default is used only for files that did not specify one.
- The number of blocks is the burden on the NameNode. Small files and small blocks both increase objects in the same direction.
- The upper limit of the number of replicas is the number of DataNodes. Even if
setrepsucceeds, check the actual number of copies withfsck. - A block is an ordinary file on a DataNode disk. It lies as a
blk_and.metapair. - The default file checksum is tied to the block layout. Compare copies with different block sizes using COMPOSITE_CRC.
What you will do in the next lab
You upload a 12MiB file with a 4MiB block and see with fsck that it splits into three blocks, then follow that block ID to find and open the actual block file on the DataNode's local disk. You also compare that uploading the same file with the default block size gives one block. You raise the replication factor to 3 and see how under-replication is reported in an environment with only one DataNode, and count that several small files become as many blocks. Finally, you confirm that the checksums of two files that differ only in block size are different under the default method and become the same under COMPOSITE_CRC.