IP Library › Granted Patent US 11,880,290
Granted Patent B2
US 11,880,290 · App. 18/165,257 · Granted Jan 23, 2024

Scalable exactly-once data processing using transactional streaming writes

Inventors: Pavan Edara (Mountain View, CA); Reuven Lax (Mountain View, CA); Yi Yang (Mountain View, CA); Gurpreet Singh Nanda (Seattle, WA)
Assignee: Google LLC
G06F11/3034G06F9/30047G06F9/467G06F11/0757G06F11/0772G06F11/1402G06F12/0246G06F12/0253G06F2201/84
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,880,290
App. No.
18/165,257
Granted
Jan 23, 2024
Kind
B2
Abstract

A method for processing data exactly once using transactional stream writes includes receiving, from a client, a batch of data blocks for storage on memory hardware in communication with the data processing hardware. The batch of data blocks is associated with a corresponding sequence number and represents a number of rows of a table stored on the memory hardware. The method also includes partitioning the batch of data blocks into a plurality of sub-batches of data blocks. For each sub-batch of data blocks, the method further includes assigning the sub-batch of data blocks to a buffered stream; writing, using the assigned buffered stream, the sub-batch of data blocks to the memory hardware; updating a storage log with an intent to commit the sub-batch of data blocks using the assigned buffered stream; and committing the sub-batch of data blocks to the memory hardware.

Claims (42)

1. A computer-implemented method executed by data processing hardware that causes the data processing hardware to perform operations comprising:

receiving, from a client, a batch of data blocks for storage on memory hardware in communication with the data processing hardware;

partitioning the batch of data blocks into a plurality of sub-batches of data blocks;

assigning each sub-batch of data blocks of the plurality of sub-batches of data blocks to a respective buffered stream, each respective buffered stream configured to write the assigned sub-batches to the memory hardware;

determining that a particular sub-batch of data blocks of the plurality of sub-batches of data blocks failed to be written to the memory hardware;

in response to determining that the particular sub-batch of data blocks of the plurality of sub-batches of data blocks failed to be written to the memory hardware, assigning the particular sub-batch of data blocks to a new respective buffered stream; and

receiving an intent to commit the particular sub-batch of data blocks using the new respective buffered stream, the intent to commit indicating that the particular sub-batch of data blocks is successfully written to the memory hardware.

2. The method of claim 1 , wherein the batch of data blocks are associated with a corresponding sequence number and representing a number of rows of a table stored on the memory hardware.

3. The method of claim 1 , wherein the operations further comprise, in response to determining that the particular sub-batch of data blocks of the plurality of sub-batches of data blocks failed to be written to the memory hardware, retrying, using the respective assigned buffered stream, to write the particular sub-batch of data blocks to the memory hardware.

4. The method of claim 3 , wherein the operations further comprise determining that retrying, using the respective assigned buffered stream, to write the particular sub-batch of data blocks to the memory hardware has failed to complete before committing the particular sub-batch of data blocks to the memory hardware.

5. The method of claim 4 , wherein the operations further comprise removing, from the memory hardware, the particular sub-batch of data blocks from the respective assigned buffered stream.

6. The method of claim 5 , wherein removing, from the memory hardware, the particular sub-batch of data blocks from the respective assigned buffered stream comprises performing garbage-collection on the particular sub-batch of data blocks from the respective assigned buffered stream.

7. The method of claim 1 , wherein the operations further comprise:

in response to receiving the intent to commit the particular sub-batch of data blocks to the memory hardware, determining a current timestamp; and

associating the particular sub-batch of data blocks with the current timestamp.

8. The method of claim 7 , wherein the operations further comprise converting the particular sub-batch of data blocks into a read-optimized format based on the current timestamp.

9. The method of claim 7 , wherein the operations further comprise:

receiving a query request at a snapshot timestamp, the query request requesting return of data blocks stored on the memory hardware that match query parameters; and

returning any data blocks of the particular sub-batch of data blocks that match the query parameters when the snapshot timestamp is later than the current timestamp associated with the particular sub-batch of data blocks.

10. The method of claim 1 , wherein the intent to commit indicates committing the particular sub-batch of data blocks to the memory hardware using a flush transform.

11. A system comprising:

data processing hardware; and

memory hardware in communication with the data processing hardware, the memory hardware storing instructions that when executed on the data processing hardware cause the data processing hardware to perform operations comprising:

receiving, from a client, a batch of data blocks for storage on memory hardware in communication with the data processing hardware;

partitioning the batch of data blocks into a plurality of sub-batches of data blocks;

assigning each sub-batch of data blocks of the plurality of sub-batches of data blocks to a respective buffered stream, each respective buffered stream configured to write the assigned sub-batches to the memory hardware;

determining that a particular sub-batch of data blocks of the plurality of sub-batches of data blocks failed to be written to the memory hardware;

in response to determining that the particular sub-batch of data blocks of the plurality of sub-batches of data blocks failed to be written to the memory hardware, assigning the particular sub-batch of data blocks to a new respective buffered stream; and

receiving an intent to commit the particular sub-batch of data blocks using the new respective buffered stream, the intent to commit indicating that the particular sub-batch of data blocks is successfully written to the memory hardware.

12. The system of claim 11 , wherein the batch of data blocks are associated with a corresponding sequence number and representing a number of rows of a table stored on the memory hardware.

13. The system of claim 11 , wherein the operations further comprise, in response to determining that the particular sub-batch of data blocks of the plurality of sub-batches of data blocks failed to be written to the memory hardware, retrying, using the respective assigned buffered stream, to write the particular sub-batch of data blocks to the memory hardware.

14. The system of claim 13 , wherein the operations further comprise determining that retrying, using the respective assigned buffered stream, to write the particular sub-batch of data blocks to the memory hardware has failed to complete before committing the particular sub-batch of data blocks to the memory hardware.

15. The system of claim 14 , wherein the operations further comprise removing, from the memory hardware, the particular sub-batch of data blocks from the respective assigned buffered stream.

16. The system of claim 15 , wherein removing, from the memory hardware, the particular sub-batch of data blocks from the respective assigned buffered stream comprises performing garbage-collection on the particular sub-batch of data blocks from the respective assigned buffered stream.

17. The system of claim 11 , wherein the operations further comprise:

in response to receiving the intent to commit the particular sub-batch of data blocks to the memory hardware, determining a current timestamp; and

associating the particular sub-batch of data blocks with the current timestamp.

18. The system of claim 17 , wherein the operations further comprise converting the particular sub-batch of data blocks into a read-optimized format based on the current timestamp.

19. The system of claim 17 , wherein the operations further comprise:

receiving a query request at a snapshot timestamp, the query request requesting return of data blocks stored on the memory hardware that match query parameters; and

returning any data blocks of the particular sub-batch of data blocks that match the query parameters when the snapshot timestamp is later than the current timestamp associated with the particular sub-batch of data blocks.

20. The system of claim 11 , wherein the intent to commit indicates committing the particular sub-batch of data blocks to the memory hardware using a flush transform.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Feb 6, 2023
From: EDARA, PAVAN; LAX, REUVEN; YANG, YI; NANDA, GURPREET SINGH
To: GOOGLE LLC
Reel/Frame 062606/0710 →
Continuity (2)
Continuation 17085576 · Oct 30, 2020
Related Publication 20230185688A1 · Jun 15, 2023