A YARN application's containers run on dozens or hundreds of NodeManagers, and each container writes its stdout, stderr and framework logs to local disk on whichever node it landed on. When the job fails at 3 a.m., those logs are the evidence, and they are scattered across machines that will delete them a few hours later. Log aggregation is the YARN feature that collects each application's container logs into one place on a shared file system (normally HDFS or an object store) so that yarn logs and the history server UI can show them after the containers, and even the nodes, are gone.

This article explains the pipeline from first principles: where container logs start, how the NodeManager uploads them, the directory layout and file formats on HDFS, rolling aggregation for long-running services, policies that choose which containers to keep, retention and the small-files cost, and how to read the result. It finishes with a worked incident, failure modes and a checklist. Property names and defaults are from yarn-default.xml in current Apache Hadoop; check your distribution's version.

The pipeline at a glance

Containerstdout, stderr, syslogNM local log diryarn.nodemanager.log-dirswritesNodeManagerper-app log aggregatorreadsHDFS remote log dir/app-logs/alice/bucket-logs/0042/application_..._0042/host1_8041 (one file per node)upload on finishResourceManageraggregation statusreportyarn logs / history UIyarn.log.server.urlreadDeletion serviceretain-secondsdelete old appsContainers write locally; the NodeManager uploads one file per node per application; readers and the deletion service work on HDFS.
Containers log locally; each NodeManager uploads one file per application to HDFS; readers and the deletion service work from there.

Where container logs start

Every container gets a log directory under one of the NodeManager's yarn.nodemanager.log-dirs (default ${yarn.log.dir}/userlogs), laid out as <log-dir>/<application_id>/<container_id>/. The launch script redirects the process's stdout and stderr there, and frameworks put their own files alongside (MapReduce's syslog, Spark's executor logs). Details of the launch are in YARN containers.

Aggregation is off by default (yarn.log-aggregation-enable=false). Without it, the NodeManager keeps an application's logs for yarn.nodemanager.log.retain-seconds (default 10,800 seconds, three hours) after it finishes and then deletes them. You can only read them through that node's web UI while it is up. For any shared cluster, turn aggregation on.

The aggregation pipeline

With aggregation enabled, each NodeManager runs a log aggregator per application that had containers on it. The sequence for a normal batch application:

  1. Containers run and write to local log directories.
  2. When the application finishes, the ResourceManager tells each NodeManager involved. Each aggregator collects the logs of that application's containers on its node.
  3. The aggregator writes a single file for the node into the application's remote directory, named after the node (for example host1.example.com_8041), first under a temporary name and then renamed into place, so readers do not see half-written files.
  4. The NodeManager reports aggregation status to the ResourceManager, which shows it on the application page as, for example, SUCCEEDED, RUNNING or FAILED. If a NodeManager does not report within yarn.log-aggregation-status.time-out.ms (default 600,000 ms), the RM marks it timed out.
  5. With yarn.log-aggregation.enable-local-cleanup true (the default), the local copies are deleted after upload.

The important property of this design is one file per node per application. An application whose containers touched 300 nodes produces 300 files, regardless of how many containers ran. That number drives both the NameNode cost and the read cost later.

The remote directory layout

The remote root is yarn.nodemanager.remote-app-log-dir (default /tmp/logs). Move it out of /tmp, which some sites clean on schedule. Current Hadoop releases write into a bucketed layout to keep directory sizes bounded: the application directory sits under the user, a bucket- prefix plus the suffix from yarn.nodemanager.remote-app-log-dir-suffix (default logs), and a four-digit bucket equal to the application's sequence number modulo 10,000.

# yarn-site.xml (excerpt)
yarn.log-aggregation-enable                 = true
yarn.nodemanager.remote-app-log-dir         = /app-logs
yarn.nodemanager.remote-app-log-dir-suffix  = logs
yarn.log-aggregation.retain-seconds         = 1209600     # 14 days; default -1 keeps forever
yarn.log.server.url = http://jhs.example.com:19888/jobhistory/logs

# Resulting path for application_1727700000000_0042 owned by alice:
/app-logs/alice/bucket-logs/0042/application_1727700000000_0042/host1.example.com_8041

# The root must exist with the sticky bit, writable by all, like /tmp:
hdfs dfs -mkdir -p /app-logs
hdfs dfs -chown yarn:hadoop /app-logs
hdfs dfs -chmod 1777 /app-logs

The NodeManager checks the root's permissions and expects 1777; it logs a warning if they differ. The sticky bit lets every user's applications create their own subdirectories without being able to delete each other's. When reading, yarn.nodemanager.remote-app-log-dir-include-older (default true) lets tools also find applications written in the older, non-bucketed layout, which matters on clusters upgraded with logs in place.

File formats: TFile and indexed

How each per-node file is written is pluggable. yarn.log-aggregation.file-formats lists controllers, and the first one is used for writing; the others remain readable. The default controller is TFile (LogAggregationTFileController), the original format: a container-by-container stream that is simple and universally supported. The Indexed format (LogAggregationIndexedFileController) adds an index so a reader can seek to one container's file without scanning the whole node file, and it was designed with rolling uploads in mind.

yarn.log-aggregation.file-formats = IFile,TFile        # write IFile, still read TFile
yarn.log-aggregation.file-controller.IFile.class =
    org.apache.hadoop.yarn.logaggregation.filecontroller.ifile.LogAggregationIndexedFileController
yarn.log-aggregation.file-controller.TFile.class =
    org.apache.hadoop.yarn.logaggregation.filecontroller.tfile.LogAggregationTFileController

Keep TFile in the list after a switch, as the default file's own description recommends, so logs already on HDFS stay readable. Check that every client that reads logs (your history server, any tooling that parses files directly) is on a release that understands the indexed format before you change the writer.

Rolling aggregation for long-running applications

Batch jobs upload once, at the end. Long-running services such as a Spark Streaming job or a service running on YARN never end, so by default their logs would never be aggregated and the local disks would fill. Rolling aggregation fixes that: yarn.nodemanager.log-aggregation.roll-monitoring-interval-seconds (default -1, off) makes NodeManagers upload periodically while the application runs. A floor, roll-monitoring-interval-seconds.min (default 3,600), stops anyone setting it so low that the NameNode is flooded. yarn.nodemanager.log-aggregation.num-log-files-per-app (default 30) caps how many aggregated files are kept per application per node, deleting the oldest.

Rolling works when the application rotates its own log files, because the NodeManager selects files to upload by name, using the rolled-logs pattern. The application declares which files are rolled through its LogAggregationContext. In Spark, spark.yarn.rolledLog.includePattern sets this for you, paired with a rolling file appender in the executor's logging configuration. A hand-written YARN client sets it directly:

// Java YARN client: aggregate everything at the end, and roll *.log.N files
// while the application is running.
LogAggregationContext ctx = LogAggregationContext.newInstance(".*", "");
ctx.setRolledLogsIncludePattern(".*\\.log\\.[0-9]+");
ctx.setLogAggregationPolicyClassName(
    "org.apache.hadoop.yarn.server.nodemanager.containermanager.logaggregation"
    + ".FailedOrKilledContainerLogAggregationPolicy");

ApplicationSubmissionContext app = yarnClient.createApplication()
    .getApplicationSubmissionContext();
app.setLogAggregationContext(ctx);

Choosing which containers to keep

Not every container's logs are worth keeping. A policy decides which containers are aggregated. The cluster default is yarn.nodemanager.log-aggregation.policy.class (default AllContainerLogAggregationPolicy), with parameters in yarn.nodemanager.log-aggregation.policy.parameters; an application can override both in its context, as above. The implementations in the NodeManager's log aggregation package:

PolicyAggregatesUse for
AllContainerLogAggregationPolicyEvery containerDefault; small clusters, audit needs
FailedOrKilledContainerLogAggregationPolicyContainers that failed or were killedLarge, mostly successful batch jobs
FailedContainerLogAggregationPolicyFailed containers onlySame, stricter
AMOrFailedContainerLogAggregationPolicyThe ApplicationMaster plus failed containersKeeping the job's narrative and its errors
AMOnlyLogAggregationPolicyThe ApplicationMaster onlyVery wide jobs where task logs are noise
SampleContainerLogAggregationPolicyAM, failed containers and a sample of successful onesBig jobs where you want some baseline
NoneContainerLogAggregationPolicyNothingWorkloads that ship logs elsewhere
LimitSizeContainerLogAggregationPolicyContainers within a size limitGuarding against log floods

The sample policy takes parameters such as SR:0.2,MIN:20 (its defaults): every container is kept while the application has at most 20, and beyond that a 20% sample of successful ones. Changing the cluster default to keep only failed and AM containers is the single largest lever on aggregated log volume, at the cost of not having a successful task's log to compare against.

Reading aggregated logs

The yarn logs command reads aggregated logs (and, for running applications, can fetch from NodeManagers). The useful options:

# What is there, without dumping it: containers, nodes, file names and sizes
yarn logs -applicationId application_1727700000000_0042 -show_application_log_info
yarn logs -applicationId application_1727700000000_0042 -show_container_log_info

# The ApplicationMaster's logs (attempt 1), only stderr
yarn logs -applicationId application_1727700000000_0042 -am 1 -log_files stderr

# One container, last 64 KB of each file (negative size reads from the end)
yarn logs -applicationId application_1727700000000_0042 \
  -containerId container_e03_1727700000000_0042_01_000007 -size -65536

# Everything to a local directory, one file per container
yarn logs -applicationId application_1727700000000_0042 -out /tmp/app42-logs

Access follows the application's ACLs and the HDFS permissions on the user's directory, so a user can read their own applications' logs and an admin can read everyone's. The web path is set by yarn.log.server.url: the RM and NodeManager UIs redirect finished applications' log links to it, normally the MapReduce JobHistory Server. For finished-application metadata beyond logs, see the Timeline Server.

Retention and NameNode cost

Aggregated logs are kept forever unless you set yarn.log-aggregation.retain-seconds (default -1). The deletion service that enforces it runs inside the MapReduce JobHistory Server, checking every retain-check-interval-seconds (default: one tenth of the retention). If you do not run a JobHistory Server, nothing deletes old logs; confirm which daemon runs it in your distribution.

The cost is NameNode objects, not bytes. Take a cluster running 3,000 applications a day, each touching 80 nodes on average: that is 240,000 files a day, plus directories. Thirty days of retention means over 7 million files held in NameNode memory, a meaningful fraction of a mid-sized NameNode's capacity (the small files problem explains the arithmetic). The levers, in order of effect: a stricter aggregation policy, shorter retention, a namespace quota on the log root so a runaway job hits a limit instead of the NameNode (see HDFS quotas), and mapred archive-logs, which packs eligible applications' log files into Hadoop archives.

Worked example: logs that stopped arriving

Worked example. Newly onboarded users report that yarn logs returns nothing for their jobs from the last two days. The RM application page shows log aggregation status FAILED for many of them. On one of those nodes, the NodeManager log shows repeated upload failures with a permission error on /app-logs.

The history: a cleanup script run by an administrator recursively reset permissions under /app-logs to 755. Users who already had a directory under the root were unaffected, but the first applications of every new user failed to create their directory, and the NodeManagers kept local copies only until yarn.nodemanager.log.retain-seconds expired. The fix was hdfs dfs -chmod 1777 /app-logs (the root only, not recursive), then confirming new applications reached SUCCEEDED aggregation. Logs from the two days were lost on nodes past the local retention window. Two follow-ups prevent a repeat: a monitoring check on the root's permission bits, and an alert on the ratio of applications whose aggregation status is not SUCCEEDED.

Failure modes

  • Aggregation never enabled. Logs vanish three hours after a job ends; discovered during the first serious incident.
  • Wrong root permissions. Uploads fail and logs are silently kept only locally, then deleted.
  • Long-running service without rolling. Local disks fill, the NodeManager's disk health checker marks directories bad, and the node stops taking containers (see the NodeManager guide).
  • Rolling without rotation. Rolling is configured but the application writes one ever-growing file, so nothing matches the rolled pattern.
  • Infinite retention. The log root grows into millions of files and becomes a NameNode heap problem.
  • Log floods. A job printing a stack trace per record writes gigabytes per container; aggregation then spends minutes uploading it and HDFS space disappears.

Trade-offs

Aggregation trades NameNode objects and HDFS space for the ability to debug after the fact. Keeping all containers is the best for forensics and the worst for the NameNode; failed-plus-AM policies cut volume dramatically but remove the healthy baseline. Rolling aggregation makes services debuggable but adds files every interval. Many large sites also ship logs to a search system, which gives full-text search across applications; YARN aggregation stays the system of record because it works without any extra infrastructure and respects application ACLs.

What to do next

  1. Enable aggregation and move the remote root out of /tmp; create it with 1777 permissions.
  2. Set yarn.log-aggregation.retain-seconds to an agreed retention, and confirm the JobHistory Server (or your distribution's equivalent) runs the deletion service.
  3. Count files under the log root, and put a namespace quota on it.
  4. Choose a cluster default policy; consider AMOrFailedContainerLogAggregationPolicy for large batch clusters.
  5. Configure rolling aggregation and file rotation for every long-running application.
  6. Point yarn.log.server.url at the history server and test a link from the RM UI for a finished job.
  7. Alert on applications whose aggregation status is not SUCCEEDED.
Key takeaway: YARN log aggregation copies each application's container logs from NodeManager disks into one file per node on HDFS, where yarn logs and the history server can read them after the containers are gone. Enable it, give the remote root 1777 permissions outside /tmp, set retention and confirm something enforces it, use rolling uploads with file rotation for long-running services, choose an aggregation policy that fits your volume, and watch both aggregation status and the file count under the log root.