IP Library › Granted Patent US 10,223,184
Granted Patent B1
US 10,223,184 · App. 14/036,792 · Granted Mar 5, 2019

Individual write quorums for a log-structured distributed storage system

Inventors: Samuel James McKelvie (Seattle, WA); Benjamin Tobler (Seattle, WA); James McClellan Corey (Bothell, WA); Pradeep Jnana Madhavarapu (Mountain View, CA); Oscar Ricardo Moll Thomae (Seattle, WA); Christopher Richard Newcombe (Kirkland, WA); Yan Valerie Leshinsky (Kirkland, WA); Anurag Windlass Gupta (Atherton, CA)
Assignee: Amazon Technologies, Inc.
G06F11/0727G06F11/0709
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 10,223,184
App. No.
14/036,792
Granted
Mar 5, 2019
Kind
B1
Abstract

A log-structured distributed storage system may implement individual write quorums. Log records may be sent to different storage nodes of a quorum set storing data for a storage client sufficient to satisfy a write quorum requirement. For each log record, acknowledgments from storage nodes are received, and a determination is made whether the write quorum requirement is satisfied for the log record. Different log records may be maintained at different storage nodes, and still satisfy the write quorum requirement such that in some embodiments no one storage node may maintain all of the log records sent to storage nodes in the quorum set.

Claims (77)

1. A system, comprising:

a plurality of storage nodes implementing a quorum set of a distributed storage system, wherein each storage node in the quorum set maintains a log-structured data store storing data for a storage client; and

a client-side storage driver module implemented on one or more computing devices of the storage client, configured to:

receive one or more updates to the data;

generate a plurality of log records indicating the one or more updates;

send each of the plurality of log records to at least some of the storage nodes of the quorum set sufficient to satisfy a write quorum requirement;

for each log record:

receive acknowledgments from at least some storage nodes of the quorum set; and

determine that the acknowledgments received for the log record satisfy the write quorum requirement indicating that the log record is made durable at the quorum set, wherein at least one of the storage nodes of the quorum set satisfying the write quorum requirement for one of the plurality of log records is different from the storage nodes of the quorum set satisfying the write quorum requirement for another of the plurality of log records;

determine that a respective instance of a sequence completion point for different ones of the storage nodes of the quorum set is the same sufficient to satisfy a recovery quorum requirement;

identify the sequence completion point as a truncation point in a log record sequence for the quorum set for a recovery operation;

recover log records in the log record sequence generated up to the truncation point, wherein log records in the log record sequence generated after the truncation point are excluded from the recovery operation; and

make the data available for processing access requests.

2. The system of claim 1 , wherein the client-side storage driver module is further configured to:

maintain storage system metadata for each log record that indicates the different storage nodes of the quorum set acknowledging the log record;

receive a read request associated with a particular view of the data;

in response to receiving the read request:

based, at least in part, on the storage system metadata, identify at least one storage node of the quorum set to service the read request that maintains the particular view of the data; and

send the read request to the at least one storage node.

3. The system of claim 1 , wherein each of the plurality of storage nodes in the quorum set is configured to:

evaluate the log records stored at the storage node according to a the log record sequence in order to determine a sequence completion point for the log records stored at the storage node;

obtain another sequence completion point for each of one or more other storage nodes in the quorum set;

based, at least in part, on the sequence completion point for the storage node and the obtained other sequence completion points for each of the one or more other storage nodes in the quorum set, identify at least one storage node of the one or more other storage nodes in the quorum set advanced further in the log record sequence than the storage node;

request one or more additional log records from the identified at least one storage node that complete the log record sequence between the sequence completion point for the storage node and the other sequence completion point for the identified at least one storage node; and

in response to receiving the requested one or more additional log records, advance the sequence completion point for the storage node to include the one or more additional log records.

4. The system of claim 1 , wherein each of the plurality of storage nodes implementing the quorum set maintains a respective sequence completion point for log records stored at the storage node according to the log record sequence, and wherein the recovery operation is an operation to recover from failure of the storage client.

5. A method, comprising:

performing, by a plurality of computing devices:

sending each of a plurality of log records indicating updates to data to at least some storage nodes of a plurality of storage nodes implementing a quorum set storing the data for a storage client that are sufficient to satisfy a write quorum requirement, wherein the plurality of log records are sent according to a log record sequence;

for each log record:

receiving acknowledgments of the log record from at least some storage nodes of the quorum set;

determining that the acknowledgments received for the log record satisfy the write quorum requirement indicating that the log record is made durable at the quorum set, wherein at least one of the storage nodes of the quorum set satisfying the write quorum requirement for one of the plurality of log records is different from the storage nodes of the quorum set satisfying the write quorum requirement for another of the plurality of log records;

determining that a respective instance of a sequence completion point for different ones of the storage nodes of the quorum set is the same sufficient to satisfy a recovery quorum requirement;

identifying the sequence completion point as a truncation point in the log record sequence for the quorum set for a recovery operation;

recovering log records in the log record sequence generated up to the truncation point, wherein log records in the log record sequence generated after the truncation point are excluded from the recovery operation; and

making the data available for processing access requests.

6. The method of claim 5 , wherein no individual storage node of the different storage nodes maintains all of the plurality of log records.

7. The method of claim 5 , wherein at least some of the plurality of log records are received at the different ones of the plurality of storage nodes in an order different than the log record sequence.

8. The method of claim 5 , further comprising evaluating, at each of one or more of the plurality of storage nodes in the quorum set, log records maintained at the storage node according to the log record sequence in order to determine a respective sequence completion point for the log records maintained at the respective storage node.

9. The method of claim 8 , further comprising:

performing, by each of the one or more storage nodes:

receiving another sequence completion point from each of one or more other storage nodes in the quorum set;

based, at least in part, on the respective sequence completion point for the storage node and the received other sequence completion point for each of the one or more other storage nodes in the quorum set, identifying at least one storage node in the quorum set advanced further in the log record sequence than the storage node;

requesting one or more additional log records from the identified at least one storage node that complete the log record sequence between the respective sequence completion point for the storage node and the received others sequence completion point for the identified at least one storage node; and

in response to receiving the requested one or more additional log records, advancing the respective sequence completion point for the storage node to include the one or more additional log records.

10. The method of claim 8 , wherein said sending is performed by the storage client, wherein said evaluating is performed by each of the one or more storage nodes of the plurality of storage nodes in the quorum set, and wherein the recovery operation is an operation to recover from failure of the storage client.

11. The method of claim 10 , wherein one or more storage nodes of the plurality of storage nodes implementing the quorum set fails such that the write quorum requirement for one or more log records maintained by the respective sequence completion point for at least one remaining storage node of the quorum set is not the same sufficient to satisfy the recovery quorum requirement, wherein said identifying the truncation point in the log record sequence, and said sending the one or more additional log records to the different ones of the plurality of storage nodes implementing the quorum are performed with regard to the at least one remaining storage node.

12. The method of claim 5 , further comprising:

maintaining storage system metadata that identifies the different ones of the storage nodes in the quorum set that acknowledge each of the plurality of log records;

receiving a read request associated with a particular view of the data;

selecting at least one storage node in the quorum set to service the read request based, at least in part, on the storage system metadata; and

sending the read request to the selected at least one storage node in order to be serviced.

13. The method of claim 5 , wherein the storage client is a database node implemented as part of a network-based database service, and wherein the plurality of storage nodes in the quorum set are implemented as part of a network-based storage service.

14. A non-transitory, computer-readable storage medium, storing program instructions that when executed by one or more computing devices cause the one or more computing devices to implement a client-side storage driver module that implements:

sending each of a plurality of log records indicating updates to data to at least some storage nodes of a plurality of storage nodes implementing a quorum set storing the data for a storage client that are sufficient to satisfy a write quorum requirement, wherein the plurality of log records are sent in order according to a log record sequence;

for each log record:

receiving acknowledgments of the log record from at least some storage nodes of the quorum set, wherein acknowledgements for at least a subset of the plurality of log records are received in an order different than an order corresponding to the log record sequence; and

determining that the acknowledgments received for the log record satisfy the write quorum requirement indicating that the log record is made durable at the quorum set, wherein at least one of the storage nodes of the quorum set satisfying the write quorum requirement for one of the plurality of log records is different from the storage nodes of the quorum set satisfying the write quorum requirement for another of the plurality of log records; and

determining that a respective instance of a sequence completion point for different ones of the storage nodes of the quorum set is the same sufficient to satisfy a recovery quorum requirement;

identifying the sequence completion point as a truncation point in the log record sequence for the quorum set for a recovery operation;

recovering log records in the log record sequence generated up to the truncation point, wherein log records in the log record sequence generated after the truncation point are excluded from the recovery operation; and

making the data available for processing access requests.

15. The non-transitory, computer-readable storage medium of claim 14 , wherein, in said sending the plurality of log records corresponding to the log record sequence to be maintained at the different ones of the plurality of storage nodes, the client-side storage driver module implements grouping different ones in the plurality of log records into one or more batches of log records to be sent to the different ones of the plurality of storage nodes together.

16. The non-transitory, computer-readable storage medium of claim 14 , wherein the client-side storage driver module further implements:

maintaining storage system metadata that identifies the different ones of the storage nodes in the quorum set that acknowledge each of the plurality of log records;

receiving a read request associated with a particular view of the data;

selecting at least one storage node in the quorum set to service the read request based, at least in part, on the storage system metadata; and

sending the read request to the selected at least one storage node in order to be serviced.

17. The non-transitory, computer-readable storage medium of claim 14 , wherein the one or more storage nodes acknowledging the log record in satisfaction of the write quorum requirement different from the storage nodes acknowledging the other log record in satisfaction of the write quorum requirement is one or more additional storage nodes added to the quorum set as a result of not receiving an acknowledgment for a particular write request within a period of time.

18. The non-transitory, computer-readable storage medium of claim 14 , wherein the program instructions cause another plurality of computing devices implementing the plurality of storage nodes in the quorum set to each implement:

evaluating log records maintained at the storage node according to the log record sequence in order to determine a respective sequence completion point for the log records maintained at the respective storage node;

receiving another sequence completion point from each of one or more other storage nodes in the quorum set;

based, at least in part, on the respective sequence completion point for the storage node and the received other sequence completion point for each of the one or more other storage nodes in the quorum set, identifying at least one storage node in the quorum set advanced further in the log record sequence than the storage node;

requesting one or more additional log records from the identified at least one storage node that complete the log record sequence between the respective sequence completion point for the storage node and the received other sequence completion point for the identified at least one storage node; and

in response to receiving the requested one or more additional log records, advancing the respective sequence completion point for the storage node to include the one or more additional log records.

19. The non-transitory, computer-readable storage medium of claim 14 , wherein the recovery operation is an operation to recover from failure of the storage client.

20. The non-transitory, computer-readable storage medium of claim 14 , wherein the client-side storage driver module is implemented on a database node that is part of a network-based distributed database service, and wherein the plurality of storage nodes are implemented as part of a network-based, multi-tenant storage service.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Jan 21, 2014
From: MCKELVIE, SAMUEL JAMES; TOBLER, BENJAMIN; COREY, JAMES MCCLELLAN; MADHAVARAPU, PRADEEP JNANA; THOMAE, OSCAR RICARDO MOLL; NEWCOMBE, CHRISTOPHER RICHARD; LESHINSKY, YAN VALERIE; GUPTA, ANURAG WINDLASS
To: AMAZON TECHNOLOGIES, INC.
Reel/Frame 032008/0194 →
Cited By (5)
US 12,216,642 US 12,375,478 US 12,517,889 US 12,572,290 US 12,572,291