IP Library › Granted Patent US 11,556,431
Granted Patent B2
US 11,556,431 · App. 17/236,505 · Granted Jan 17, 2023

Rollback recovery with data lineage capture for data pipelines

Inventors: Eric Simon (Juvisy, FR); Cesar Salgado Vieira de Souza (Porto Alegre, BR)
Assignee: SAP SE
G06F11/1469G06F12/0253G06F2201/84G06F2212/702
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,556,431
App. No.
17/236,505
Granted
Jan 17, 2023
Kind
B2
Abstract

Computer-readable media, methods, and systems are disclosed for performing rollback recovery with data lineage capture for data pipelines. A middle operator receives ingested input events from a source operator reading data from an external input data source. The middle operator then logs information regarding middle input events to a middle operator input log, designating the logged middle input event information as incomplete. The middle operator then processes data associated with the middle input events and updates the middle input log entries setting them to a completed logging status designation for middle input events that were consumed to produce the one or more middle output events. The middle operator then transmits the middle output events to subsequent operators. Garbage collection is performed to remove completed entries from the middle operator output log. Finally, based on receiving a recovering message from a subsequent operator, corresponding middle output events are re-sent.

Claims (52)

1. One or more non-transitory computer-readable media storing computer-executable instructions that, when executed by a processor, perform a method for performing rollback recovery with data lineage capture for data pipelines, the method comprising:

at a middle operator, receiving, from a source operator, one or more input events ingested by the source operator by way of a read operation to an external input data source;

logging information regarding one or more middle input events to a middle operator input log associated with the middle operator, wherein the one or more middle input events are logged with an incomplete logging status designation;

processing data associated with the one or more middle input events;

updating one or more middle input log entries by setting the one or more middle input log entries to a completed logging status designation corresponding to a consumed subset of the one or more middle input events that were consumed to produce one or more middle output events;

transmitting the one or more middle output events to one or more subsequent operators; and

based on receiving a recovering message from one or more subsequent operators, resending corresponding middle output events from a middle operator output log.

2. The non-transitory computer-readable media of claim 1 , the method further comprising:

performing background garbage collection on the middle operator output log, wherein updated middle input log events that have been updated to reflect the completed logging status designation are removed from the middle operator output log.

3. The non-transitory computer-readable media of claim 1 , the method further comprising:

establishing one or more data lineage analysis start points and one or more data lineage analysis target points.

4. The non-transitory computer-readable media of claim 3 , wherein the updating one or more middle input log entries by setting the one or more middle input log entries to a completed logging status designation comprises setting the one or more middle input log entries to a completed log preservation status designation.

5. The non-transitory computer-readable media of claim 2 , wherein the performing background garbage collection on the middle operator output log comprises preserving the one or more middle input log entries having a completed log preservation status designation.

6. The non-transitory computer-readable media of claim 3 , the method further comprising:

inserting a monitoring agent operator into the data pipeline to separate output events of data lineage interest from a remainder of output events; and

establishing the one or more data lineage analysis start points downstream from the inserted monitoring agent operator.

7. The non-transitory computer-readable media of claim 6 , the method further comprising:

traversing the middle operator input log and the middle operator output log to identify intermediate input and output values at intermediate operators between the one or more data lineage analysis start points and the one or more data lineage analysis target points to determine initial values of the one or more data lineage analysis target points.

8. A method for performing rollback recovery with data lineage capture for data pipelines, the method comprising:

at a middle operator, receiving, from a source operator, one or more input events ingested by the source operator by way of a read operation to an external input data source;

logging information regarding one or more middle input events to a middle operator input log associated with the middle operator, wherein the one or more middle input events are logged with an incomplete logging status designation;

processing data associated with the one or more middle input events;

updating one or more middle input log entries by setting the one or more middle input log entries to a completed logging status designation corresponding to a consumed subset of the one or more middle input events that were consumed to produce one or more middle output events;

transmitting the one or more middle output events to one or more subsequent operators; and

based on receiving a recovering message from one or more subsequent operators, resending corresponding middle output events from a middle operator output log.

9. The method of claim 8 , further comprising:

performing background garbage collection on the middle operator output log, wherein updated middle input log events that have been updated to reflect the completed logging status designation are removed from the middle operator output log.

10. The method of claim 8 , the method further comprising:

establishing one or more data lineage analysis start points and one or more data lineage analysis target points.

11. The method of claim 10 , wherein the updating one or more middle input log entries by setting the one or more middle input log entries to a completed logging status designation comprises setting the one or more middle input log entries to a completed log preservation status designation.

12. The method of claim 9 , wherein the performing background garbage collection on the middle operator output log comprises preserving the one or more middle input log entries having a completed log preservation status designation.

13. The method of claim 10 , the method further comprising:

inserting a monitoring agent operator into the data pipeline to separate output events of data lineage interest from a remainder of output events; and

establishing the one or more data lineage analysis start points downstream from the inserted monitoring agent operator.

14. The method of claim 13 , the method further comprising:

traversing the middle operator input log and the middle operator output log to identify intermediate input and output values at intermediate operators between the one or more data lineage analysis start points and the one or more data lineage analysis target points to determine initial values of the one or more data lineage analysis target points.

15. A system comprising at least one processor and at least one non-transitory memory storing computer executable instructions that when executed by the processor cause the system to carry out actions comprising:

at a middle operator, receiving, from a source operator, one or more input events ingested by the source operator by way of a read operation to an external input data source;

logging information regarding one or more middle input events to a middle operator input log associated with the middle operator, wherein the one or more middle input events are logged with an incomplete logging status designation;

processing data associated with the one or more middle input events;

updating one or more middle input log entries by setting the one or more middle input log entries to a completed logging status designation corresponding to a consumed subset of the one or more middle input events that were consumed to produce one or more middle output events;

transmitting the one or more middle output events to one or more subsequent operators;

performing background garbage collection on a middle operator output log, wherein updated middle input log events that have been updated to reflect the completed logging status designation are removed from the middle operator output log; and

based on receiving a recovering message from one or more subsequent operators, resending corresponding middle output events from the middle operator output log.

16. The system of claim 15 , wherein the middle operator input log and the middle operator output log are combined to form a middle operator merged log.

17. The system of claim 15 , the actions further comprising:

establishing one or more data lineage analysis start points and one or more data lineage analysis target points.

18. The system of claim 17 , wherein the updating one or more middle input log entries by setting the one or more middle input log entries to a completed logging status designation comprises setting the one or more middle input log entries to a completed log preservation status designation.

19. The system of claim 18 , wherein the performing background garbage collection on the middle operator output log comprises preserving the one or more middle input log entries having a completed log preservation status designation.

20. The system of claim 17 , the actions further comprising:

inserting a monitoring agent operator into a data pipeline to separate output events of data lineage interest from a remainder of output events; and

establishing the one or more data lineage analysis start points downstream from the inserted monitoring agent operator.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 21, 2021
From: SIMON, ERIC; VIEIRA DE SOUZA, CESAR SALGADO
To: SAP SE
Reel/Frame 055991/0645 →
Continuity (1)
Related Publication 20220365851A1 · Nov 17, 2022