skip to content

In Hadoop MapReduce, what is an InputSplit and what determines its size?

level: middleimportance: should knowfreq 42%

answer

  1. logical work, not stored bytes
  2. one of these, one map task
  3. a formula with three inputs
  4. the block size is only the default
  5. a gzip file has exactly one

basics

~20 s

An InputSplit is the logical byte range one map task will process, plus host hints for locality — not a physical copy of data. With FileInputFormat its size is max(minSize, min(maxSize, blockSize)), so by default it equals the HDFS block size of 128 MB.

solid answer

~40 s

An **InputSplit** is a logical unit of work: a file, an offset, a length, and the hosts holding that data. The framework launches exactly one map task per split, so splits determine map parallelism. `FileInputFormat` computes the size as `max(minSize, min(maxSize, blockSize))`, where `minSize` comes from `mapreduce.input.fileinputformat.split.minsize` (1 by default), `maxSize` from `mapreduce.input.fileinputformat.split.maxsize` (unset), and `blockSize` from `dfs.blocksize` (128 MB in Hadoop 3). So by default one split ≈ one block, which is what gives you data locality. Two things break the mapping: a **non-splittable** codec such as gzip forces one split per file no matter how large, and **many small files** produce one tiny split each, flooding the cluster with short map tasks — `CombineFileInputFormat` or compaction is the fix. Splits are logical, so a record straddling a boundary is still read whole.

code

text · 9 lines
text
dfs.blocksize                                  = 134217728   (128 MB)
mapreduce.input.fileinputformat.split.minsize  = 1
mapreduce.input.fileinputformat.split.maxsize  = unset

splitSize = max(minsize, min(maxsize, blocksize)) = 134217728

1 GB uncompressed text file  ->  8 splits  ->  8 map tasks
1 GB single gzip file        ->  1 split   ->  1 map task (codec not splittable)
5000 files of 2 MB each      -> 5000 splits -> 5000 map tasks (one per file)

go deeper

for a junior

Recall that one map task processes one InputSplit and that a split defaults to the HDFS block size, so map parallelism roughly tracks input size divided by 128 MB.

for a middle

State the formula max(minSize, min(maxSize, blockSize)) with the config keys behind each term, and explain the logical-versus-physical distinction including how a record that crosses a boundary is still read whole.

for a senior

Diagnose the failure modes in production: a non-splittable gzip input serializing an entire job, and a small-files directory generating tens of thousands of short tasks. Know both the job-side workaround and the upstream file-layout fix.

for a principal

Own the ingestion contract that prevents these problems: target file sizes, splittable container formats with internal compression, and compaction policy. Split behaviour is downstream of how you let data land, and that is a platform decision.

## Logical work, not physical storage An InputSplit is a description of work, produced by the `InputFormat` before any task starts: typically a path, a start offset, a length in bytes, and a list of hosts that hold those bytes. It contains no data. The framework creates exactly one map task per split and passes the split to a `RecordReader`, which is what actually opens the stream and produces key/value pairs. Keeping splits logical is what lets one `InputFormat` sit over HDFS files, object storage, a database table, or a Hive partition without changing the execution model. An HDFS block, by contrast, is physical: a fixed-size chunk of a file stored and replicated on DataNodes. The two coincide by default and are constantly confused in interviews. The split's host list is derived from where the block replicas live, which lets YARN try to place the map container on a node that already holds the data — node-local, then rack-local, then anywhere. ## The size formula `FileInputFormat.computeSplitSize` is `max(minSize, min(maxSize, blockSize))` where: - `mapreduce.input.fileinputformat.split.minsize` — default 1 byte - `mapreduce.input.fileinputformat.split.maxsize` — unset, i.e. effectively unbounded - `dfs.blocksize` — 128 MB by default in Hadoop 2 and 3 With the defaults, the formula collapses to the block size. To get **more, smaller** map tasks, lower `split.maxsize`; to get **fewer, larger** ones, raise `split.minsize` above the block size — which deliberately sacrifices some data locality, since a split then spans blocks that may live on different nodes. Note what you cannot do: `mapreduce.job.maps` is advisory and is ignored by file-based input formats. The number of mappers is a consequence of the split calculation, never a number you set directly. One more detail worth knowing: `FileInputFormat` will not create a tiny trailing split. It applies a slop factor of 1.1, so a remainder no bigger than 10% of a split size is folded into the previous split rather than getting a map task of its own. ## Records that cross a boundary If splits are byte ranges and records are lines, a line will eventually straddle a boundary. `LineRecordReader` handles it with a simple convention: a split whose start offset is not zero skips forward past the first newline (that partial line belongs to the previous split), and every reader keeps reading *past* its own end until it finishes the record in progress. That final read may pull a few bytes from a block on another node — a small, bounded remote read. The consequence is that records are never split or duplicated, but it also means custom `InputFormat`s over record structures without a cheap delimiter must implement this logic themselves or be marked non-splittable. ## Compression changes everything Splittability is a property of the codec, not of the file size: - **gzip** — not splittable. A 10 GB `.gz` file yields one split and therefore one map task that reads all 10 GB, while the rest of the cluster idles. This is one of the most common causes of an inexplicably serial "parallel" job. - **bzip2** — splittable, at a notable CPU cost for compression and decompression. - **LZO** — splittable only if an index file has been generated for it. - **Snappy** — not splittable on its own for raw text, but perfectly fine *inside* a block-compressed container such as SequenceFile, Avro, ORC or Parquet, where the container defines the split points. The practical rule: store large inputs in a container format with internal block compression, or keep individual raw files near the block size. ## The small-files problem, from the job's side Because a split never spans two files in the default `FileInputFormat`, a directory of 50,000 files of 2 MB each produces 50,000 splits and 50,000 map tasks. Each is a container allocation, a JVM start, and a few seconds of real work — the scheduling overhead dwarfs the computation. `CombineFileInputFormat` packs many small files (preferring ones on the same node and rack) into a single split with a target size, which collapses that to a manageable task count. The durable fix is upstream: compact small files before processing. ## What to say in an interview Lead with "logical range, one map task per split, defaults to the block size", then show the formula and immediately name the two things that break the naive picture — non-splittable compression and small files. Those two are what an interviewer is actually probing for, because they are the causes of most surprising map-side parallelism in real clusters.

  • How do you increase the number of map tasks for a large uncompressed input file?
    Lower `mapreduce.input.fileinputformat.split.maxsize` below the block size; the split formula then produces more, smaller splits and therefore more map tasks. Setting `mapreduce.job.maps` does nothing for file-based inputs — it is advisory and ignored. The tradeoff is more container allocations and JVM starts, so pushing splits far below the block size usually costs more than it saves.
  • Why can a 10 GB gzip file make a cluster look completely idle during a job?
    gzip is not splittable, so `FileInputFormat` emits a single split for the whole file and the job runs one map task reading all 10 GB while every other node sits idle. Fixes are to store the data in a splittable container (ORC, Parquet, Avro, SequenceFile) or in bzip2 or indexed LZO, or to break the input into files near the block size.
  • What does CombineFileInputFormat change about split creation?
    It allows a single split to span multiple files, packing small files up to a target size and preferring files that share a node and rack so locality is largely preserved. A directory of tens of thousands of tiny files then produces a few hundred splits instead of tens of thousands of short map tasks. It treats the symptom; compacting the files upstream is the real fix.

saying these in an interview costs you the question

  • Says an InputSplit is a physical chunk of the file on disk
  • Claims you set the mapper count with mapreduce.job.maps
  • Assumes any compressed input is split like an uncompressed one
  • Believes a record crossing a split boundary is dropped or duplicated
  • Thinks split size is always exactly the HDFS block size

context