Many MapReduce jobs need the same side data in every task: a lookup table of product names for a log enricher, a stop-word list, a trained model, a native library or a Python virtual environment for a Streaming job. Reading it from HDFS inside every task works, but hundreds of tasks hammering the same few blocks at startup is slow and puts load on a handful of DataNodes. The distributed cache solves this. You declare files and archives when you submit the job; YARN copies each one to every node that runs a task, once, unpacks archives, and links them into each container's working directory, where the task opens them as ordinary local files.
This article explains the mechanism end to end: the API and command-line options, how NodeManagers localize and share files, the visibility rules that decide whether a copy is reused across jobs, a worked map-side join, sizing, failure modes and alternatives.
Not to be confused with HDFS caching
The name causes confusion, so first, what it is not. HDFS also has centralized cache management, which pins HDFS blocks in DataNode memory so reads of hot data avoid disk; that is a storage-layer feature covered in HDFS caching. The MapReduce distributed cache is a job-submission feature: it ships read-only files to the local disk of compute nodes for the life of a job, and it has nothing to do with memory caching. The two are independent and can be combined.
How a cached file reaches a task
At submission, the job client records each cache resource's URI together with its modification time, size and visibility in the job configuration. If you point at a local file (through -files), the client first copies it into the job's staging directory on the cluster file system. When the ApplicationMaster asks a NodeManager to launch a container, the launch request lists these as local resources.
The NodeManager's ResourceLocalizationService then downloads each resource that is not already cached on that node, verifies that its timestamp still matches, unpacks archives, and creates a symlink in the container's working directory named after the URI fragment (or the file name). Only after every resource is localized does the container process start. Containers on the same node share one copy, and with the right visibility, so do later jobs. See the NodeManager article for where localization sits in the container lifecycle.
The API and command-line options
The Job class has four methods. addCacheFile(URI) ships a file. addCacheArchive(URI) ships an archive that the node unpacks; zip, jar, tar, tar.gz and tgz are recognized, and the symlink then points at the unpacked directory. addFileToClassPath(Path) and addArchiveToClassPath(Path) ship a jar or archive and add it to the task classpath. The old static DistributedCache class is deprecated; use the Job methods. A #name fragment on the URI sets the symlink name, which lets you give files stable names regardless of the HDFS path:
Job job = Job.getInstance(conf, "enrich-clicks");
job.setJarByClass(EnrichDriver.class);
// Side data, linked as ./products.tsv in every task's working directory.
job.addCacheFile(new URI("hdfs:///ref/products/2026-10-02/products.tsv#products.tsv"));
// A whole directory tree, unpacked and linked as ./geo/
job.addCacheArchive(new URI("hdfs:///ref/geo/geo-2026q3.tgz#geo"));
// An extra dependency on the task classpath.
job.addFileToClassPath(new Path("/libs/fastutil-8.5.13.jar"));From the command line, drivers that run through ToolRunner get the generic options -files, -archives and -libjars, which take comma-separated local or remote paths and accept the same #name fragments. Hadoop Streaming uses them to ship scripts and environments:
hadoop jar enrich.jar com.example.EnrichDriver \
-files hdfs:///ref/products/2026-10-02/products.tsv#products.tsv \
-archives hdfs:///envs/py311-enrich.tgz#venv \
-libjars /opt/libs/fastutil-8.5.13.jar \
/data/clicks/2026-10-02 /out/enriched/2026-10-02
mapred streaming -files mapper.py,stopwords.txt \
-mapper "venv/bin/python mapper.py" -reducer NONE \
-input /data/docs -output /out/tokensGeneric options must come before the job's own arguments, and they only work if the driver implements Tool; a driver that builds its Configuration by hand silently ignores them.
Worked example: a map-side join
The classic use is a map-side (replicated) join: one large dataset, one small one that fits in task memory. Instead of shuffling both to reducers, every mapper loads the small side once in setup() and joins each record as it streams past:
public class EnrichMapper extends Mapper<LongWritable, Text, Text, Text> {
private final Map<String, String> products = new HashMap<>();
@Override
protected void setup(Context ctx) throws IOException {
// The symlink created from the #products.tsv fragment.
try (BufferedReader r = Files.newBufferedReader(Paths.get("products.tsv"))) {
String line;
while ((line = r.readLine()) != null) {
int tab = line.indexOf('\t');
products.put(line.substring(0, tab), line.substring(tab + 1));
}
}
ctx.getCounter("enrich", "products_loaded").increment(products.size());
}
@Override
protected void map(LongWritable k, Text v, Context ctx) throws IOException, InterruptedException {
String[] f = v.toString().split("\t", -1); // click: ts, user, sku, ...
String name = products.get(f[2]);
if (name == null) { ctx.getCounter("enrich", "unknown_sku").increment(1); return; }
ctx.write(new Text(f[2]), new Text(f[0] + "\t" + f[1] + "\t" + name));
}
}Open the file by its symlink name, relative to the working directory. context.getCacheFiles() returns the original URIs, which is useful for logging which version a task loaded, but not a path to open. The job can be map-only, with no shuffle at all; compare the reduce-side join in shuffle and sort, which moves both datasets across the network.
Visibility: who shares a localized copy
Whether a localized copy is shared depends on its visibility, which the client decides at submission. A PUBLIC resource is shared by all users and jobs on the node. A PRIVATE resource is shared only among jobs of the same user. An APPLICATION resource belongs to one job and is deleted when the job finishes; the job jar, configuration and split files are APPLICATION resources. Local files passed with -files or -libjars are uploaded to the job's private staging directory, so they are PRIVATE and, because every job gets a new staging path, they are localized afresh on every run.
For files you add yourself, the client marks a resource PUBLIC only if the file is world-readable and every ancestor directory up to the root is world-executable; otherwise it is PRIVATE. This matters in practice: a reference file under a user home directory with mode 700 is PRIVATE, so it is downloaded again for each user who runs the job. If many users run jobs against the same reference data, put it under a directory such as /ref with 755 directories and 644 files, so every NodeManager downloads it once for everyone. Do not do this with sensitive data: PUBLIC copies are readable from any job on the node.
Lifecycle and retention on each node
PUBLIC and PRIVATE copies outlive the job. Each NodeManager keeps them in its local directories and evicts least-recently-used resources that no running container is using when the total exceeds a target size. The relevant settings, from yarn-default.xml:
| Property | Default | Meaning |
|---|---|---|
yarn.nodemanager.localizer.cache.target-size-mb | 10240 | Per-node retention target for PUBLIC and PRIVATE resources (APPLICATION excluded) |
yarn.nodemanager.localizer.cache.cleanup.interval-ms | 600000 | How often cleanup runs (10 minutes) |
yarn.nodemanager.localizer.fetch.thread-count | 4 | Parallel downloads for public resources |
yarn.nodemanager.local-dirs | under hadoop.tmp.dir | Where localized files live; spread across disks |
The target is a retention goal, not a hard cap: resources in use are never deleted, so a node running jobs with large archives can exceed it. Size the local directories for peak working set plus intermediate map output, which shares the same disks; YARN containers covers the rest of a container's local footprint.
Sizing: network and memory
Worked example: a 300 MB product file, a job of 2,000 map tasks, a cluster of 50 nodes. Reading from HDFS in every task would move 2,000 x 300 MB, about 600 GB, mostly from the three DataNodes holding the file's replicas. With the distributed cache, each node downloads once: 50 x 300 MB, about 15 GB, spread over time as containers start. Run the job daily with PUBLIC visibility and an unchanged file, and later runs download nothing.
Memory is the other half. A 300 MB text file becomes considerably more as a Java HashMap of String objects, often two to four times its size because of object headers, references and UTF-16 strings. With a 2 GB map heap that is uncomfortably tight. Shrink the data before shipping it: project only the needed columns, encode keys as longs, and use primitive collections. Measure the loaded size once in a single test task with a heap dump or the GC log rather than guessing. If the small side will not fit, it is not small: use a reduce-side join.
Debugging localization
When a task cannot find its file, look at what was actually localized. A task can log its working directory listing in setup(); symlinks there show which cached path each name resolved to. On the node, localized files sit under the NodeManager local directories in filecache (PUBLIC), usercache/<user>/filecache (PRIVATE) and usercache/<user>/appcache/<application id> (per-application data and container working directories). YARN deletes container directories as soon as the container exits; on a test cluster, set yarn.nodemanager.delete.debug-delay-sec (default 0) to a few minutes so you can inspect them, and never leave it high in production. Localization failures appear in the NodeManager log and in the container diagnostics shown by yarn application -status and the job history page, usually naming the resource and the reason.
Failure modes
- File changed after submission. Overwriting the HDFS file while a job runs makes localization fail on nodes that have not fetched it yet, with an error that the resource changed on the source file system. Fix: never overwrite; publish each version to a new, dated path and point jobs at it.
- Wrong path in the task. Opening the HDFS path, or an absolute local path, instead of the symlink name. Fix: always use a
#namefragment and open that relative name. - Ignored generic options.
-filespassed to a driver that does not useToolRunner, or placed after positional arguments. The job runs, and every task fails with file-not-found. Fix: implementTool. - Out of memory in setup(). The side table grew. Fix: a counter of loaded entries, a size check in the driver that refuses to submit above a threshold, and the shrinking tactics above.
- Local disk exhaustion. Large archives retained by many concurrent jobs fill the NodeManager disks and the node is marked unhealthy. Fix: lower the target size, watch local directory usage, and keep archives lean.
- Slow container start. A multi-gigabyte environment archive localized on cold nodes delays every first task. Fix: PUBLIC visibility so later runs reuse the copy, and smaller environments.
Alternatives and trade-offs
The distributed cache is the right tool for read-only side data from kilobytes to a few hundred megabytes that most tasks need. For data each task needs only a slice of, read that slice from HDFS directly. For dependencies, a fat job jar avoids per-jar localization but makes every build larger. On Spark, the equivalents are --files and --archives for files and broadcast variables for in-memory lookup tables, and the memory arithmetic above applies unchanged. For the map-side join pattern in the context of mapper design, see writing mappers.
Two refinements are worth knowing. Hadoop 3 includes an optional YARN shared cache (designed in YARN-1492): a cluster service that keeps checksummed copies of job resources such as jars so that repeated submissions skip the upload to the staging directory. It needs its own manager service and is off by default, so check your distribution's documentation before relying on it. Second, version your side data the way you version code: a driver that resolves a dated path such as /ref/products/2026-10-02/ and logs it gives every output a traceable input, and makes a rerun reproducible even after newer reference data lands.
What to do next
- Find jobs that read the same HDFS file in every task (look for side-file paths in
setup()), and move that file to the distributed cache. - Publish reference data to immutable, dated paths, and give each cache entry a
#namefragment. - Check permissions on shared reference directories so resources localize as PUBLIC.
- Make sure every driver implements
Tooland runs throughToolRunner. - Measure the heap cost of each side table in one task, and add a driver-side size guard.
- Monitor NodeManager local-directory usage and container localization time, and tune the cache target size if nodes run short of disk.