IP Library › Granted Patent US 11,500,852
Granted Patent B2
US 11,500,852 · App. 16/914,075 · Granted Nov 15, 2022

Database system with database engine and separate distributed storage service

Inventors: Anurag Windlass Gupta (Atherton, CA); Neal Fachan (Seattle, WA); Samuel James McKelvie (Seattle, WA); Laurion Darrell Burchall (Seattle, WA); Christopher Richard Newcombe (Kirkland, WA); Pradeep Jnana Madhavarapu (Mountain View, CA); Benjamin Tobler (Seattle, WA); James McClellan Corey (Bothell, WA)
Assignee: Amazon Technologies, Inc.
G06F16/2365G06F11/1451G06F11/1471G06F16/23G06F11/2094G06F2201/80
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 11,500,852
App. No.
16/914,075
Granted
Nov 15, 2022
Kind
B2
Abstract

A database system may include a database service and a separate distributed storage service. The database service (or a database engine head node thereof) may be responsible for query parsing, optimization, and execution, transactionality, and consistency, while the storage service may be responsible for generating data pages from redo log records and for durability of those data pages. For example, in response to a write request directed to a particular data page, the database engine head node may generate a redo log record and send it, but not the data page, to a storage service node. The storage service node may store the redo log record and return a write acknowledgement to the database service prior to applying the redo log record. The server node may apply the redo log record and other redo log records to a previously stored version of the data page to create a current version.

Claims (52)

1. A system, comprising:

at least one processor; and

a memory, storing program instructions that when executed by the at least one processor cause the at least one processor to implement a database system, the database system configured to:

receive a request to add to a database table or modify the database table, wherein the request is received from a client of the database system;

in response to the reception of the request, send a log record that describes the addition to the database table or the modification to the database table to a plurality of storage nodes of a distributed storage system that store data of the database table to which the log record is subsequently applied; and

after receiving write acknowledgements for the request to add to the database table or modify the database table from a number of the plurality of storage nodes that satisfies a write quorum, return an acknowledgement of the request as complete to the client.

2. The system of claim 1 , wherein the plurality of storage nodes is a protection group for a portion of a database volume that includes the data.

3. The system of claim 1 , wherein the database system is further configured to:

receive a query that targets the database table;

as part of performing the query, send respective requests for the data to the plurality of storage nodes;

after receiving the data from a number of the storage nodes that satisfy a read quorum, use the data to perform the query; and

return a result of the query.

4. The system of claim 3 , wherein the log record is applied to update the data at the plurality of storage nodes after the respective requests for the data are received at the storage nodes.

5. The system of claim 3 , wherein the respective requests for the data are sent to the plurality of storage nodes after determining that a cache at the database system does not store the data.

6. The system of claim 1 , wherein the database system is further configured to:

receive a different request to add to the database table or modify the database table, wherein the request is received from the client of the database system;

send a different one or more log records that describe the different request to the plurality of storage nodes; and

after a period of time without receiving enough write acknowledgements from a number of the plurality of storage nodes that satisfy the write quorum, determine that an error condition exists in the distributed storage system.

7. The system of claim 1 , wherein the database system is a database web service and wherein the distributed storage system is a separate, log-structured storage web service and wherein the log record is written to respective append-only storage at the plurality of storage nodes.

8. A method, comprising:

receiving, at a database system, a request to add to a database table or modify the database table, wherein the request is received from a client of the database system;

in response to the receiving of the request, sending, by the database system, a log record that describes the addition to the database table or the modification to the database table to a plurality of storage nodes of a distributed storage system that store data of the database table to which the log record is subsequently applied; and

after receiving write acknowledgements for the request to add to the database table or modify the database table from a number of the plurality of storage nodes that satisfies a write quorum, returning, by the database system, an acknowledgement of the request as complete to the client.

9. The method of claim 8 , wherein the plurality of storage nodes is a protection group for a portion of a database volume that includes the data.

10. The method of claim 8 , further comprising:

receiving a query that targets the database table;

as part of performing the query, sending respective requests for the data to the plurality of storage nodes;

after receiving the data from a number of the storage nodes that satisfy a read quorum, using the data to perform the query; and

returning a result of the query.

11. The method of claim 10 , wherein the log record is applied to update the data at the plurality of storage nodes after the respective requests for the data are received at the storage nodes.

12. The method of claim 10 , wherein the respective requests for the data are sent to the plurality of storage nodes after determining that a cache at the database system does not store the data.

13. The method of claim 8 , further comprising:

receiving a different request to add to the database table or modify the database table, wherein the request is received from the client of the database system;

sending a different one or more log records that describe the different request to the plurality of storage nodes; and

after a period of time without receiving enough write acknowledgements from a number of the plurality of storage nodes that satisfy the write quorum, determining that an error condition exists in the distributed storage system.

14. The method of claim 8 , wherein the database system is a database web service and wherein the distributed storage system is a separate, log-structured storage web service and wherein the log record is written to respective append-only storage at the plurality of storage nodes.

15. One or more non-transitory, computer-readable storage media, storing program instructions that when executed on or across one or more computing devices, cause the one or more computing devices to implement:

receiving, at a database system, a request to add to a database table or modify the database table, wherein the request is received from a client of the database system;

in response to the receiving of the request, sending, by the database system, a log record that describes the addition to the database table or the modification to the database table to a plurality of storage nodes of a distributed storage system that store data of the database table to which the log record is subsequently applied; and

after receiving write acknowledgements for the request to add to the database table or modify the database table from a number of the plurality of storage nodes that satisfies a write quorum, returning, by the database system, an acknowledgement of the request as complete to the client.

16. The one or more non-transitory, computer-readable storage media of claim 15 , wherein the plurality of storage nodes is a protection group for a portion of a database volume that includes the data.

17. The one or more non-transitory, computer-readable storage media of claim 15 , storing further instructions that when executed on or across the one or more computing devices, cause the one or more computing devices to further implement:

receiving a query that targets the database table;

as part of performing the query, sending respective requests for the data to the plurality of storage nodes;

after receiving the data from a number of the storage nodes that satisfy a read quorum, using the data to perform the query; and

returning a result of the query.

18. The one or more non-transitory, computer-readable storage media of claim 17 , wherein the log record is applied to update the data at the plurality of storage nodes after the respective requests for the data are received at the storage nodes.

19. The one or more non-transitory, computer-readable storage media of claim 17 , wherein the respective requests for the data are sent to the plurality of storage nodes after determining that a cache at the database system does not store the data.

20. The one or more non-transitory, computer-readable storage media of claim 15 , storing further instructions that when executed on or across the one or more computing devices, cause the one or more computing devices to further implement:

receiving a different request to add to the database table or modify the database table, wherein the request is received from the client of the database system;

sending a different one or more log records that describe the different request to the plurality of storage nodes; and

after a period of time without receiving enough write acknowledgements from a number of the plurality of storage nodes that satisfy the write quorum, determining that an error condition exists in the distributed storage system.

Continuity (4)
Continuation 15369681 · Dec 5, 2016
Division 14201493 · Mar 7, 2014
Provisional Application 61794572 · Mar 15, 2013
Related Publication 20200327114A1 · Oct 15, 2020
Cited By (1)
US 12,517,889