Configuring the ResourceManager for Work-Preserving Restart
Work-preserving ResourceManager restart ensures that applications continuously function during a ResourceManager restart with minimal impact to end-users. The overall concept is that the ResourceManager preserves application queue state in a pluggable state store, and reloads that state on restart. While the ResourceManager is down, ApplicationMasters and NodeManagers continuously poll the ResourceManager until it restarts. When the ResourceManager comes back online, the ApplicationMasters and NodeManagers re-register with the newly started ResourceManger. When the ResourceManager restarts, it also recovers container information by absorbing the container statuses sent from all NodeManagers. Thus, no work will be lost due to a ResourceManager crash-reboot event
To configure work-preserving restart for the ResourceManager, set the following properties in the yarn-site.xml file.
Property:
yarn.resourcemanager.recovery.enabled
Value:
true
Description:
Enables ResourceManager restart. The default value is false. If this
configuration property is set to true, running applications will resume
when the ResourceManager is restarted.
Example:
<property>
<name>yarn.resourcemanager.recovery.enabled</name>
<value>true</value>
</property>Property:
yarn.resourcemanager.store.class
Value:
<specified_state_store>
Description:
Specifies the state-store used to store application and application-attempt state and other credential information to enable restart. The available state-store implementations are:
org.apache.hadoop.yarn.server.resourcemanager.recovery.FileSystemRMStateStore– a state-store implementation persisting state to a file system such as HDFS. This is the default value.org.apache.hadoop.yarn.server.resourcemanager.recovery.LeveldbRMStateStore– a LevelDB-based state-store implementation.org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore– a ZooKeeper-based state-store implementation.
Example:
<property>
<name>yarn.resourcemanager.store.class</name>
<value>org.apache.hadoop.yarn.server.resourcemanager.recovery.FileSystemRMStateStore</value>
</property>FileSystemRMStateStore Configuration
The following properties apply only if org.apache.hadoop.yarn.server.resourcemanager.recovery.FileSystemRMStateStore has been specified as the state-store in the yarn.resourcemanager.store.class property.
Property:
yarn.resourcemanager.fs.state-store.uri
Value:
<hadoop.tmp.dir>/yarn/system/rmstore
Description:
The URI pointing to the location of the file system path where the RM state will be stored
(e.g. hdfs://localhost:9000/rmstore). The default value is
<hadoop.tmp.dir>/yarn/system/rmstore.
Example:
<property>
<name>yarn.resourcemanager.fs.state-store.uri</name>
<value>hdfs://localhost:9000/rmstore</value>
</propertyProperty:
yarn.resourcemanager.fs.state-store.retry-policy-spec
Value:
2000, 500
Description:
The Hadoop FileSystem client retry policy specification. Hadoop FileSystem client retry is
always enabled. This is specified in pairs of sleep-time and number-of-retries i.e. (t0, n0),
(t1, n1), ..., the first n0 retries sleep t0 milliseconds on average, the following n1 retries
sleep t1 milliseconds on average, and so on. The default value is (2000,
500).
Example:
<property>
<name>yarn.resourcemanager.fs.state-store.retry-policy-spec</name>
<value>2000, 500</value>
</propertyLeveldbRMStateStore Configuration
The following properties apply only if org.apache.hadoop.yarn.server.resourcemanager.recovery.LeveldbRMStateStore has been specified as the state-store in the yarn.resourcemanager.store.class property.
Property:
yarn.resourcemanager.leveldb-state-store.path
Value:
<hadoop.tmp.dir>/yarn/system/rmstore
Description:
The local path where the RM state will be stored.
Example:
<property>
<name>yarn.resourcemanager.leveldb-state-store.path</name>
<value><hadoop.tmp.dir>/yarn/system/rmstore</value>
</propertyZKRMStateStore Configuration
The following properties apply only if org.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore has been specified as the state-store in the yarn.resourcemanager.store.class property.
Property:
yarn.resourcemanager.zk-address
Value:
<host>:<port>
Description:
A comma-separated list of <host>:<port> pairs, each
corresponding to a server in a ZooKeeper cluster where the ResourceManager state will be
stored.
Example:
<property>
<name>yarn.resourcemanager.zk-address</name>
<value>127.0.0.1:2181</value>
</propertyProperty:
yarn.resourcemanager.zk-state-store.parent-path
Value:
/rmstore
Description:
The full path of the root znode where RM state will be stored. The default value is
/rmstore.
Example:
<property>
<name>yarn.resourcemanager.zk-state-store.parent-path</name>
<value>/rmstore</value>
</propertyProperty:
yarn.resourcemanager.zk-num-retries
Value:
500
Description:
The number of times the ZooKeeper-client running inside the ZKRMStateStore tries to
connect to ZooKeeper in case of connection timeouts. The default value is
500.
Example:
<property>
<name>yarn.resourcemanager.zk-num-retries</name>
<value>500</value>
</propertyProperty:
yarn.resourcemanager.zk-retry-interval-ms
Value:
2000
Description:
The interval in milliseconds between retries when connecting to a ZooKeeper server. The default value is 2 seconds.
Example:
<property>
<name>yarn.resourcemanager.zk-retry-interval-ms</name>
<value>2000</value>
</propertyProperty:
yarn.resourcemanager.zk-timeout-ms
Value:
10000
Description:
The ZooKeeper session timeout in milliseconds. This configuration is used by the ZooKeeper server to determine when the session expires. Session expiration happens when the server does not hear from the client (i.e. no heartbeat) within the session timeout period specified by this property. The default value is 10 seconds.
Example:
<property>
<name>yarn.resourcemanager.zk-timeout-ms</name>
<value>10000</value>
</propertyProperty:
yarn.resourcemanager.zk-acl
Value:
world:anyone:rwcda
Description:
The ACLs to be used for setting permissions on ZooKeeper znodes. The default value is
world:anyone:rwcda. Example
<property>
<name>yarn.resourcemanager.zk-acl</name>
<value>world:anyone:rwcda</value>
</property>
