Role of Journal nodes in Namenode HA

What is the Role of Journal nodes in Namenode HA ?

I know many of us are aware that Role of Journal nodes is to keep both the Namenodes in sync and avoid hdfs split brain scenario by allowing only Active NN to write into journals. Have you ever wonder how does it works? Here you go!

journal nodes


Journal nodes are distributed system to store edits. Active Namenode as a client writes edits to journal nodes and commit only when its replicated to all the journal nodes in a distributed system. Standby NN need to read data from edits to be in sync with Active one. It can read from any of the replica stored on journal nodes.

ZKFC will make sure that only one Namenode should be active at a time. However, when a failover occurs, it is still possible that the previous Active NameNode could serve read requests to clients, which may be out of date until that NameNode shuts down when trying to write to the JournalNodes. For this reason, we should configure fencing methods even when using the Quorum Journal Manager.


How quorum journal manager work with fencing ?

To work with fencing journal manager uses epoc numbers. Epoc numbers are integer which always gets increased and have unique value once assigned. Namenode generate epoc number using simple algorithm and uses it while sending RPC requests to the QJM. When you configure Namenode HA, the first Active Namenode will get epoc value 1. In case of  failover or restart, epoc number will get increased. The Namenode with higher epoc number is considered as newer than any Namenode with earlier epoc number.


Now both Namenode thinks that they are active and sends write request to quorum journal manager with their epoc number, how QJM handles this situation?

Quorum journal manager stores epoc number locally which called as promised epoc. Whenever JournalNode receives RPC request along with epoc number from Namenode, it compares the epoch number with promised epoch. If request is coming from newer node which means epoc number is greater than promised epoc then it records new epoc number as promised epoc. If the request is coming from Namenode with older epoc number, then QJM simply rejects the request.


When QJM rejects the requests from Namenode with older epoc value then you get below lines in the Namenode logs

WARN client.QuorumJournalManager ( - Remote journal <journal-node-hostname>:<port> failed to write txns 2397121201-2397121201. Will try to write to this JN again after the next log roll. 
org.apache.hadoop.ipc.RemoteException( IPC's epoch 112 is less than the last promised epoch 113


Please comment if you have any feedback/questions/suggestions. Happy Hadooping!! :)



facebooktwittergoogle_plusredditpinterestlinkedinmailby feather

1 comment

  • Sathish


    Is there any possibility to configure HDFS Federation with High Availability (Active/Standby). If it is, could you please post the steps.

    Thanks in advance.

    Sathish G.

Leave a Reply to Sathish Cancel reply

Your email address will not be published. Required fields are marked *

You may use these HTML tags and attributes: <a href="" title=""> <abbr title=""> <acronym title=""> <b> <blockquote cite=""> <cite> <code> <del datetime=""> <em> <i> <q cite=""> <s> <strike> <strong>