IP Library Granted Patent US 9,753,954
Granted Patent B2
US 9,753,954 · App. 14/024,585 · Granted Sep 5, 2017

Data node fencing in a distributed file system

View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 9,753,954
App. No.
14/024,585
Granted
Sep 5, 2017
Kind
B2
Abstract

Systems and methods for data node fencing in a distributed file system to prevent data inconsistencies and corruptions are disclosed. An embodiment includes implementing a protocol whereby data nodes detect a failover and determine an active name node based on transaction identifiers associated with transaction requests. The data nodes also provide to the active name node block location information and an acknowledgment. The embodiment further includes a protocol whereby a name node refrains from issuing invalidation requests to the data nodes until the name node receives acknowledgments from all data nodes that are functional.

Claims (40)

1. A method for maintaining data correctness in a Hadoop™ based distributed cluster during a failover, in which an original name node is switched to a backup name node due to failure of the original name node, the distributed cluster having a plurality of data nodes and one or more processors, the method being performed by the one or more processors and comprising:

on the backup name node:

assuming an active role to become a new active name node, upon detecting that the original name node has failed;

flagging all of the plurality of data nodes as untrusted;

for each data node among the plurality of data nodes:

queuing, instead of issuing, commands intended for a data node until the data node is flagged as trusted, and

upon receiving an acknowledgement from the data node acknowledging the assumption of the active role of the backup name node, flagging the data node as trusted; and

on a respective data node:

receiving a first command with a first transaction number from a first name node;

receiving a second command with a second transaction number from a second name node, wherein the second transaction number is greater than the first transaction number; and

sending an acknowledgment of an active role to the second name node.

2. The method of claim 1 , further comprising:

sending a message to the data node, wherein the message includes a most recent transaction identifier known to the backup name node assuming the active role.

3. The method of claim 1 , wherein commands on any block with replicated data on untrusted data nodes are queued.

4. The method of claim 1 , further comprising receiving a data report in addition to the acknowledgment of the active role from the data node.

5. The method of claim 4 , wherein the data report includes information regarding location of replicated data stored in the data node.

6. The method of claim 4 , wherein each data report includes a list of pending deletions.

7. A Hadoop™ based distributed cluster comprising an original name node, a backup name node, and a distributed file system having a plurality of data nodes,

wherein one or more processors of the backup name node are configured to perform:

assuming an active role to become a new active name node, upon detecting that the original name node has failed;

flagging all of the plurality of data nodes as untrusted;

for each data node among the plurality of data nodes:

queuing, instead of issuing, commands intended for a data node until the data node is flagged as trusted; and

upon receiving an acknowledgement from the data node acknowledging the assumption of the active role of the backup name node, flagging the data node as trusted, and

wherein one or more processors of a respective data node are configured to perform;

receiving a first command with a first transaction number from a first name node;

receiving a second command with a second transaction number from a second name node, wherein the second transaction number is greater than the first transaction number; and

sending an acknowledgment of an active role to the second name node.

8. A machine-readable storage medium having stored thereon instructions which, when executed by one or more processors, configure the processors to performs a method in a Hadoop™ based distributed cluster comprising a plurality of name nodes and a plurality of data nodes and having a distributed file system, the method comprising:

on the backup name node:

assuming an active role to become a new active name node, upon detecting that the original name node has failed;

flagging all of the plurality of data nodes as untrusted;

for each data node among the plurality of data nodes;

queuing, instead of issuing, commands intended for a data node until the data node is flagged as trusted, and

upon receiving an acknowledgement from the data node acknowledging the assumption of the active role of the backup name node, flagging the data node as trusted; and

on a respective data node:

receiving a first command with a first transaction number from a first name node;

receiving a second command with a second transaction number from a second name node, wherein the second transaction number is greater than the first transaction number; and

sending an acknowledgment of an active role to the second name node.

9. The cluster of claim 7 , wherein the data nodes are configured to ignore commands from other name nodes that issue commands having a transaction identifier lower than a transaction identifier associated with a command issued by the backup name node.

Assignments (5)
RELEASE OF SECURITY INTERESTS IN PATENTS Recorded Oct 14, 2021
From: CITIBANK, N.A.
To: CLOUDERA, INC.; HORTONWORKS, INC.
Reel/Frame 057804/0355 →
FIRST LIEN NOTICE AND CONFIRMATION OF GRANT OF SECURITY INTEREST IN PATENTS Recorded Oct 12, 2021
From: CLOUDERA, INC.; HORTONWORKS, INC.
To: JPMORGAN CHASE BANK, N.A.
Reel/Frame 057776/0185 →
SECOND LIEN NOTICE AND CONFIRMATION OF GRANT OF SECURITY INTEREST IN PATENTS Recorded Oct 12, 2021
From: CLOUDERA, INC.; HORTONWORKS, INC.
To: JPMORGAN CHASE BANK, N.A.
Reel/Frame 057776/0284 →
SECURITY INTEREST Recorded Dec 22, 2020
From: CLOUDERA, INC.; HORTONWORKS, INC.
To: CITIBANK, N.A., AS COLLATERAL AGENT
Reel/Frame 054832/0559 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Sep 17, 2013
From: LIPCON, TODD; MYERS, AARON T.; COLLINS, ELI
To: CLOUDERA, INC.
Reel/Frame 031225/0925 →