IP Library Granted Patent US 12,061,608
Granted Patent B2
US 12,061,608 · App. 16/881,631 · Granted Aug 13, 2024

Duplicate detection and replay to ensure exactly-once delivery in a streaming pipeline

Inventors: Michael Pippin (Sunnyvale, CA); David Willcox (Urbana, IL); Allie K. Watfa (Urbana, IL); George Aleksandrovich (Hoffman Estates, IL)
Assignee: YAHOO ASSETS LLC
G06F16/24556G06F9/542
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 12,061,608
App. No.
16/881,631
Granted
Aug 13, 2024
Kind
B2
Abstract

Disclosed are embodiments for providing batch performance using a stream processor. In one embodiment, a method is disclosed comprising processing a plurality of events using a stream processor and executing a deduplication process on the plurality of events using the stream processor. The plurality of events is outputted to a streaming queue and a close of books (COB) of a data transport is detected. Then, an audit process is initiated in response to detecting the COB signal, the audit process comprising comparing a set of raw events to a set of events in the streaming queue to identify a set of missing events, and replaying a set of missing events through the stream processor.

Claims (52)

1. A method comprising:

processing a plurality of events using a stream processor;

executing a deduplication process on the plurality of events using the stream processor;

outputting the plurality of events to a streaming queue; and

initiating an audit process in response to detecting a close of books (COB) signal comprising a predetermined time interval, the audit process comprising:

comparing a set of raw events to a set of events in the streaming queue to identify a set of missing events, and

replaying a set of missing events through the stream processor, the set of missing events configured to bypass the deduplication process while being re-processed.

2. The method of claim 1 , the executing the deduplication process comprising:

querying a data store using an event in the plurality of events;

dropping the event if the event exists in the data store; and

writing the event to the data store if the event does not exist in the data store.

3. The method of claim 2 , the querying the data store comprising querying a distributed database.

4. The method of claim 1 , the replaying the set of missing events comprising annotating the missing events with deduplication keys.

5. The method of claim 4 , the executing the deduplication process comprising dropping an event in the set of missing events based on a corresponding deduplication key.

6. The method of claim 1 , the replaying the set of missing events comprising:

incrementing an audit replay level; and

annotating the set of missing events using the audit replay level.

7. The method of claim 6 , the replaying the set of missing events further comprising adjusting a filter of deduplication process based on the audit replay level.

8. The method of claim 7 , the executing the deduplication process comprising dropping an event using the filter if the event is annotated with a replay level below the audit replay level.

9. A non-transitory computer-readable storage medium for tangibly storing computer program instructions capable of being executed by a computer processor, the computer program instructions defining the steps of:

processing a plurality of events using a stream processor;

executing a deduplication process on the plurality of events using the stream processor;

outputting the plurality of events to a streaming queue; and

initiating an audit process in response to detecting a close of books (COB) signal comprising a predetermined time interval, the audit process comprising:

comparing a set of raw events to a set of events in the streaming queue to identify a set of missing events, and

replaying a set of missing events through the stream processor, the set of missing events configured to bypass the deduplication process while being re-processed.

10. The computer-readable storage medium of claim 9 , the executing the deduplication process comprising:

querying a data store using an event in the plurality of events;

dropping the event if the event exists in the data store; and

writing the event to the data store if the event does not exist in the data store.

11. The computer-readable storage medium of claim 10 , the querying the data store comprising querying a distributed database.

12. The computer-readable storage medium of claim 9 , the replaying the set of missing events comprising annotating the missing events with deduplication keys.

13. The computer-readable storage medium of claim 12 , the executing the deduplication process comprising dropping an event in the set of missing events based on a corresponding deduplication key.

14. The computer-readable storage medium of claim 9 , the replaying the set of missing events comprising:

incrementing an audit replay level; and

annotating the set of missing events using the audit replay level.

15. The computer-readable storage medium of claim 14 , the replaying the set of missing events further comprising adjusting a filter of deduplication process based on the audit replay level.

16. The computer-readable storage medium of claim 15 , the executing the deduplication process comprising dropping an event using the filter if the event is annotated with a replay level below the audit replay level.

17. A system comprising:

a stream processor for processing a plurality of events, the stream processor further comprising a dedupe bolt for executing a deduplication process on the plurality of events;

a streaming queue to receive the plurality of events processed by the stream processor; and

an auditor configured to initiate an audit process in response to detecting a close of books (COB) signal comprising a predetermined time interval, the audit process comprising:

comparing a set of raw events to a set of events in the streaming queue to identify a set of missing events, and

replaying a set of missing events through the stream processor, the set of missing events configured to bypass the deduplication process while being re-processed.

18. The system of claim 17 , the executing the deduplication process comprising:

querying a data store using an event in the plurality of events;

dropping the event if the event exists in the data store; and

writing the event to the data store if the event does not exist in the data store.

19. The system of claim 17 , the replaying the set of missing events comprising annotating the missing events with deduplication keys.

20. The system of claim 17 , the replaying the set of missing events comprising:

incrementing an audit replay level; and

annotating the set of missing events using the audit replay level.

Assignments (4)
PATENT SECURITY AGREEMENT (FIRST LIEN) Recorded Sep 29, 2022
From: YAHOO ASSETS LLC
To: ROYAL BANK OF CANADA, AS COLLATERAL AGENT
Reel/Frame 061571/0773 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Dec 16, 2021
From: YAHOO AD TECH LLC (FORMERLY VERIZON MEDIA INC.)
To: YAHOO ASSETS LLC
Reel/Frame 058982/0282 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Oct 26, 2020
From: OATH INC.
To: VERIZON MEDIA INC.
Reel/Frame 054258/0635 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded May 22, 2020
From: PIPPIN, MICHAEL; WILLCOX, DAVID; WATFA, ALLIE K.; ALEKSANDROVICH, GEORGE
To: OATH INC.
Reel/Frame 052735/0271 →
Continuity (1)
Related Publication 20210365460A1 · Nov 25, 2021