YARN federation overview
Overview
YARN federation is an architectural solution that allows multiple YARN subclusters to operate as a single logical resource management system. Each subcluster manages its own compute resources, scheduling, and application lifecycle independently, while the federation provides a single submission endpoint for users and applications. Applications submitted to a YARN federation interact with the federated environment as if it is a single large YARN cluster.
YARN federation is suitable for large-scale deployments that require higher scheduling throughput, flexible workload placement, and simplified management across multiple YARN clusters. This article describes the major YARN federation components, its workflow, and concepts.
|
NOTE
YARN and HDFS federation are used together for building a full-fledged ADH federation.
|
YARN federation benefits
A federated YARN architecture provides the following benefits:
-
Horizontal scalability. YARN federation distributes resource management across multiple ResourceManagers, eliminating the scalability limitations of a single YARN cluster. In a federated YARN environment, additional compute power can be added by joining new subclusters to the federation.
-
Higher scheduling throughput. Each ResourceManager of a federation schedules applications independently within its own subcluster. By distributing scheduling across multiple ResourceManagers, a federation can handle a larger number of concurrent applications.
-
Unified application access. Applications are submitted through a single endpoint (YARN Federation Router) rather than being directly submitted to a ResourceManager. YARN Federation Router selects an appropriate subcluster and forwards all subsequent requests, such as status queries and application termination, to the correct ResourceManager.
-
Error isolation. Each subcluster operates independently. If a YARN subcluster becomes unavailable, the remaining subclusters in the pool continue to accept and execute applications, minimizing the impact of failures.
-
Flexible workload placement. Federation policies determine where applications are executed based on configurable criteria, such as cluster load, administrator-defined weights, or data locality. This helps balance workloads and optimize resource utilization across subclusters.
YARN federation components
A high-level YARN federation diagram is presented below.
The major YARN federation components are as follows.
YARN subclusters
A YARN subcluster is a regular YARN cluster that participates in a YARN federation. Each subcluster has its own ResourceManager, NodeManagers, scheduler, and compute resources. Clients and applications communicate with sublucsters through special federation components — YARN Federation Router.
When an application is submitted to a federated YARN, it selects a subcluster to run the application. This is called a home subcluster, and others are called secondary subclusters. By default, YARN attempts to complete the application, utilizing compute resources of the home subcluster. However, if required, YARN can request resources from secondary subclusters according to the federation policy.
The following diagram demonstrates basic interaction steps between YARN subclusters in a typical workflow.
YARN Federation Router
A YARN service component that acts as a client-facing entry point into the federation. In YARN federation, YARN clients communicate with YARN Federation Router using the same protocol as they would with a regular ResourceManager. The router then forwards the request to the appropriate ResourceManager within the federation. The "appropriate" ResourceManager is determined based on the information from the State store.
The main functions of YARN Federation Router:
-
Receives YARN application submission requests from clients.
-
Selects the most appropriate subcluster for running a YARN application according to a federation policy.
-
Forwards client requests (application submission, status queries, application termination, and so on) to the correct ResourceManager within a federation.
AMRMProxy
An internal YARN component that runs inside every NodeManager and acts as a proxy between an ApplicationMaster and a ResourceManager. In a federated YARN, ApplicationMaster processes communicate through the AMRMProxy instead of interacting with a ResourceManager directly.
The main AMRMProxy functions:
-
Delivers resource requests across subclusters.
-
Enforces application quotas.
-
Applies load-balancing policies.
AMRMProxy exposes the same ApplicationMasterService protocol as ResourceManager does, so from a YARN application’s perspective, it is communicating with a regular ResourceManager.
State store
A distributed storage that stores federation metadata, required for proper operation of the entire YARN federation, namely:
-
Subcluster membership. Each ResourceManager registers its subcluster within the State store and periodically sends heartbeats. This allows federation components to determine which subclusters are currently available.
-
Home subcluster mapping. A mapping that associates each YARN application with its home subcluster. This information allows YARN Federation Routers to determine where a specific application is running and where to route corresponding requests.
The State store supports a pluggable backend. In ADH, ZooKeeper is used as the default backend system for the State store data.
Federation policy store
Stores federation policies used to determine how applications and resource requests are distributed across YARN subclusters.
YARN federation policies
YARN federation policies control how applications and resource requests are distributed among YARN subclusters. YARN federation uses two types of policies:
-
Router policies. Used for selecting a home subcluster for newly submitted YARN applications.
-
AMRMProxy policies. Determine how an application’s resource requests are distributed among the home and secondary YARN subclusters.
In ADH, YARN supports weight-based policy implementation which allows routing requests to subclusters based on a special weight configuration file.
{
"routerPolicyWeights": {
"entry": [
{
"key": {
"id": "${root-cluster-id}"
},
"value": "0.5"
},{
"key": {
"id": "${sub-cluster-id}"
},
"value": "0.5"
}
]
},
"amrmPolicyWeights": {
"entry": [
{
"key": {
"id": "${root-cluster-id}"
},
"value": "0.5"
},{
"key": {
"id": "${sub-cluster-id}"
},
"value": "0.5"
}
]
},
"headroomAlpha": "1.0"
}
Enable YARN federation
You can create a YARN federation via ADCM. This operation assumes the following major steps:
-
For each YARN subcluster, set
yarn.federation.enabled=truein ADCM (Services → YARN → Primary configuration → Federation). -
For each YARN subcluster, specify a policy weights configuration (Services → YARN → Primary configuration → Federation → yarn.federation.policy-manager-params).
-
Restart the subclusters.
-
Import subcluster configurations to each other. For this, go to the Import page and select ADH clusters to import as shown in the image.
Cluster import -
Restart the subclusters.
In a federated YARN, subclusters should point to the same ZooKeeper server.
-
Copy the ZooKeeper address from your root cluster (Services → Core configuration → Primary configuration → core-site.xml → hadoop.zk.address).
-
In every subcluster, specify this ZooKeeper address in Services → Core configuration → Primary configuration → Custom core-site.xml → hadoop.zk.address.
-
In every subcluster, specify the
yarn.resourcemanager.zk-state-store.parent-pathproperty in Services → YARN → Custom yarn-site.xml.This property specifies a root znode, where ResourceManager state will be stored. The property should be set explicitly, otherwise the default znode information will be overwritten. For example:
yarn.resourcemanager.zk-state-store.parent-path=/yarn/myfed/resourcemanager -
Restart the subclusters.
For more information about YARN federation, see YARN Federation.