Every MapReduce job with more than one reducer has to answer the same question for every record a mapper emits: which reducer gets this? The answer comes from the partitioner, a single small method called once per output record. It is easy to ignore because the default works for most jobs, and it is responsible for two of the worst things a job can do: produce wrong answers with no error, and spend hours waiting on one reducer while the rest of the cluster idles.

This article treats the partitioner as a contract you have to honour. It shows where the partitioner runs, what the default does arithmetically, why some keys break it, and how to write custom partitioners for secondary sort, globally sorted output and skewed data, with code against the org.apache.hadoop.mapreduce API. Buffers, spills and merges are covered in MapReduce shuffle and sort; this page is about the decision made before any of that happens.

Advertisement

Where the partitioner runs

The partitioner decides, per record, which reducer will receive itMap task 1emits (k, v)Map task 2emits (k, v)Map task 3emits (k, v)getPartition(k, v, R)same answer in every JVMreturns p in [0, R)record stored with psorted by (p, key)partition 0partition 1partition 2hot key: 40% of datapartition 3Reducer 0Reducer 1Reducer 2the stragglerReducer 3Each reducer fetches its partition from every map output and merges by key.If two map JVMs disagree about a key's partition, that key is reduced twice, in two places, with no error.
Map outputs are partitioned per record; reducer i fetches partition i from every map task.

When a mapper calls context.write(key, value), the map task asks the partitioner for a partition number and stores the record in its in-memory buffer tagged with that number. When the buffer spills, records are sorted by partition first and by key second, so each spill file, and the final merged map output, is a sequence of contiguous partitions with an index recording where each starts. Reducer i later fetches partition i from every completed map task and merges those runs by key. The mechanics of the buffer are in the map phase and the fetch and merge in the reduce phase.

Two edge cases follow from this design. With zero reducers the job is map-only: there is no shuffle and the partitioner is never consulted. With exactly one reducer, the new-API map task does not call your partitioner at all; every record goes to partition 0. A custom partitioner that seems to work in a single-reducer test may therefore never have run.

The contract

The class is small:

public abstract class Partitioner<KEY, VALUE> {
  public abstract int getPartition(KEY key, VALUE value, int numPartitions);
}

Everything important is in the rules the signature does not state.

  1. Range. Return a value in [0, numPartitions). Anything else makes the map task fail with an IOException reporting an illegal partition, and the task attempt dies.
  2. Determinism across processes. The same key must map to the same partition in every map task, on every node, in every JVM and on every retry. Map tasks run in separate JVMs, often on different machines, and the shuffle only works if they agree.
  3. Consistency with grouping. Every key that the reducer's grouping comparator treats as equal must land in the same partition. If two keys belong to one reduce call but hash to different partitions, the group is split across two reducers.
  4. Cost. It runs once per output record, billions of times in a large job. Avoid allocation, regular expressions, string formatting and lookups in remote systems.

The value parameter exists, but partitioning on it is almost always a mistake: records with the same key and different values end up in different reducers, which breaks rule 3 for any normal grouping.

Advertisement

HashPartitioner and its arithmetic

The default partitioner is a single line:

public int getPartition(K key, V value, int numReduceTasks) {
  return (key.hashCode() & Integer.MAX_VALUE) % numReduceTasks;
}

The mask clears the sign bit so the modulo is never negative. Math.abs would look equivalent but is not, because Math.abs(Integer.MIN_VALUE) is still negative. If you write your own hash-based partitioner, copy the mask, not the abs.

The default is only as good as the key's hashCode, and three kinds of key undermine it.

  • Identity hash codes. A key class that does not override hashCode inherits Object.hashCode, which differs between JVMs. So does a Java enum, whose hashCode is final and identity-based, and so does a Java array. Each map task then sends the same logical key to a different reducer. The job succeeds, and the output contains the same key several times with partial aggregates.
  • Hash inconsistent with comparison. A key that compares case-insensitively but hashes case-sensitively sends Apple and apple to different reducers even though the sort treats them as one key. The hash must be computed from exactly the fields, normalised in exactly the way, that the grouping comparator uses.
  • Structured hash values. IntWritable and LongWritable hash to their value or a fold of it. If every key is a multiple of 10 and you run 20 reducers, only the reducers whose index is a multiple of 10 receive data. Choose a reducer count with no common factor with patterns in your keys, or mix the bits before taking the modulo.

Text hashes its bytes with a fixed function, so it is stable across JVMs; if in doubt, build composite keys from Text and primitive Writables, and test the hash, not just the output.

A custom partitioner, configurable at runtime

A typical custom partitioner routes by part of a composite key. This one partitions click events by tenant so that each reducer writes whole tenants, and reads a list of large tenants from the job configuration to give each of them a dedicated reducer. Implementing Configurable makes the framework call setConf after construction.

public class TenantPartitioner extends Partitioner<TenantDayKey, Writable>
    implements Configurable {
  private Configuration conf;
  private final Map<String, Integer> dedicated = new HashMap<>();

  @Override public void setConf(Configuration conf) {
    this.conf = conf;
    String[] big = conf.getStrings("acme.partitioner.dedicated.tenants", new String[0]);
    for (int i = 0; i < big.length; i++) dedicated.put(big[i], i);
  }
  @Override public Configuration getConf() { return conf; }

  @Override
  public int getPartition(TenantDayKey key, Writable value, int n) {
    String tenant = key.getTenant().toString();   // one allocation per record; acceptable here
    Integer slot = dedicated.get(tenant);
    int reserved = Math.min(dedicated.size(), n - 1);
    if (slot != null && slot < reserved) return slot;          // big tenants: own reducer
    int h = tenant.hashCode() & Integer.MAX_VALUE;             // String hash is stable
    return reserved + h % (n - reserved);                      // everyone else shares the rest
  }
}

// driver
job.setPartitionerClass(TenantPartitioner.class);
job.setNumReduceTasks(64);
job.getConfiguration().setStrings("acme.partitioner.dedicated.tenants", "t-0042", "t-0917");

The partitioner depends only on the tenant field, so it agrees with any grouping that groups by tenant or by tenant and day. It is safe for every value of n, including when there are more dedicated tenants than reducers. The property name is this job's own invention, not a Hadoop key.

Secondary sort: partition by part of the key

MapReduce sorts by key but not by value, so when a reducer needs each group's values in order, for example a user's events by time, the order has to move into the key. The pattern needs a composite key and three cooperating classes: the key, here (userId, timestamp); a partitioner that uses only userId; a sort comparator that orders by userId then timestamp; and a grouping comparator that compares only userId, so a single reduce call sees all of a user's records, already in time order.

public class UserPartitioner extends Partitioner<UserTimeKey, Text> {
  public int getPartition(UserTimeKey k, Text v, int n) {
    return (k.getUserId().hashCode() & Integer.MAX_VALUE) % n;   // natural key only
  }
}
public class UserGrouping extends WritableComparator {
  protected UserGrouping() { super(UserTimeKey.class, true); }
  public int compare(WritableComparable a, WritableComparable b) {
    return ((UserTimeKey) a).getUserId().compareTo(((UserTimeKey) b).getUserId());
  }
}

job.setPartitionerClass(UserPartitioner.class);
job.setSortComparatorClass(UserTimeComparator.class);      // userId, then timestamp
job.setGroupingComparatorClass(UserGrouping.class);

Forgetting the partitioner is the classic bug. With the default HashPartitioner, the hash covers the whole composite key including the timestamp, so a user's records spread across every reducer and each reducer sees fragments of every user's timeline.

Globally sorted output: TotalOrderPartitioner

Hash partitioning sorts within each reducer's output but not across them. To make the concatenation of part files globally sorted, partition by range: reducer 0 gets the smallest keys, reducer R-1 the largest. TotalOrderPartitioner does this from a file of R-1 split points, usually produced by sampling the input with InputSampler.

job.setNumReduceTasks(100);
job.setPartitionerClass(TotalOrderPartitioner.class);
Path partitionFile = new Path("/tmp/sort-job/_partitions");   // on HDFS
TotalOrderPartitioner.setPartitionFile(job.getConfiguration(), partitionFile);

// sample up to 10,000 keys, each with probability 0.01, from at most 50 splits
InputSampler.Sampler<Text, Text> sampler =
    new InputSampler.RandomSampler<>(0.01, 10000, 50);
InputSampler.writePartitionFile(job, sampler);              // writes 99 split points

Every task reads the file itself, from the path's own filesystem; only when the path is left at the default does the partitioner read it from the task's local working directory, which is the older pattern of shipping the file through the distributed cache under that name. It parses split points as the map output key class and checks that it holds exactly numReduceTasks minus one keys, failing with a wrong-number-of-partitions error otherwise, so the reducer count must be set before sampling and must not change afterwards. For keys that implement BinaryComparable, such as Text, it builds a trie over the split points by default; otherwise it binary-searches them with the job's sort comparator. The path is configurable with mapreduce.totalorderpartitioner.path and defaults to _partition.lst.

Two limits matter. The sampler draws keys from the job's input format, so this works directly only when the map output keys are the input keys, as in a sort of a SequenceFile; otherwise you need a preliminary job that writes the real output keys. And range partitioning cannot split a single key: if one key holds 30 percent of the data, the reducer that owns it gets at least 30 percent, whatever the split points say.

KeyFieldBasedPartitioner for text keys

Hadoop Streaming jobs, and Java jobs with delimited Text keys, can partition on a subset of fields without writing a class. KeyFieldBasedPartitioner takes sort-style field specifications.

mapred streaming \
  -D stream.map.output.field.separator=. \
  -D stream.num.map.output.key.fields=3 \
  -D map.output.key.field.separator=. \
  -D mapreduce.partition.keypartitioner.options=-k1,2 \
  -D mapreduce.job.reduces=32 \
  -partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner \
  -input /logs/in -output /logs/out -mapper ./map.py -reducer ./reduce.py

The first two properties tell streaming that the first three dot-separated fields of each mapper output line form the key; the next two tell the partitioner to split that key on dots and hash fields one and two, so all records sharing those two fields meet in one reducer while still sorting on all three.

Skew: the arithmetic and the fixes

Even a correct partitioner cannot balance data the key distribution makes unbalanced. Suppose a job shuffles 2 TB to 200 reducers, about 10 GB each on average, and one key carries 8 percent of the records. Its reducer receives at least 160 GB, sixteen times the average; if the others finish in 6 minutes, that reducer takes over an hour and holds the job. Adding reducers does nothing for it, because one key always lands in exactly one partition. The reduce-input counters in the shuffle article show this shape directly.

  • Combine first. For associative, commutative aggregations, a combiner shrinks a hot key's records before the shuffle, often enough on its own.
  • Salt hot keys. Append a salt in [0, s) to known hot keys so they spread across s partitions, aggregate partially, then strip the salt and finish in a second, much smaller job. This works for sums, counts and maxima, not for logic that needs every record of a key together.
  • Dedicated partitions. Route known giants to their own reducers, as the tenant example does, so they at least do not share a reducer with other work.
  • Fix the model. Skew often comes from a default value such as an empty string or a NULL user id. Filter or route those records separately.

The same reasoning applies to Spark's hash partitioning and range partitioning; see Spark shuffle.

Failure modes

SymptomCauseFix
Same key appears in several part files with partial totalsIdentity hashCode: enum, array or class without hashCodeHash a stable representation; add a test that runs the partitioner in two JVMs
Groups split across reducers in secondary sortDefault partitioner hashing the full composite keyPartition on the natural key only
Task fails with an illegal partition IOExceptionNegative or out-of-range return valueMask the sign bit; handle every value of numPartitions
Custom partitioner never seems to runJob has one reducer, or is map-onlyTest with at least two reducers
Wrong number of partitions in keysetReducer count changed after samplingSet numReduceTasks before writePartitionFile
One reducer runs for hoursHot key or a default-value keyCombiner, salting, dedicated partitions, filter defaults
Half the reducers receive nothingKey hash pattern shares a factor with the reducer countChange the reducer count or mix the hash

Trade-offs

ChoiceGainsCosts
HashPartitionerNo setup; balanced for well-distributed keysNo global order; at the mercy of hashCode
TotalOrderPartitionerGlobally sorted output; adapts to the key distributionSampling pass, a partition file to manage, cannot split one key
Custom routingBusiness-aware placement, dedicated slots for giantsCode to test and maintain; easy to violate the grouping rule
SaltingSpreads hot keys across reducersA second aggregation stage; only for decomposable logic

What to do next

  1. For every job with custom keys, confirm that hashCode is overridden, stable across JVMs and computed from the same fields as the grouping comparator.
  2. Write a unit test that calls your partitioner over a sample of real keys with several reducer counts and asserts the range and the distribution.
  3. Check every secondary-sort job for the trio: partitioner on the natural key, full sort comparator, grouping comparator on the natural key.
  4. Read reduce input record counters for your slowest jobs; a top reducer far above the median is skew, not a hardware problem.
  5. Add a combiner where the aggregation allows it before reaching for salting or custom routing.
  6. Use TotalOrderPartitioner only when downstream consumers need globally sorted files, and fix the reducer count before sampling.
Key takeaway: The partitioner is one method with three obligations: return a value in range, return the same value for the same key in every JVM, and agree with the reducer's grouping. Most partitioner bugs are hashCode bugs that produce plausible wrong output rather than errors. Use HashPartitioner when keys hash well, partition on the natural key for secondary sort, use TotalOrderPartitioner with sampling for global order, and treat skew as a property of your keys that only combining, salting or remodelling can fix.