IP Library › Granted Patent US 11,711,255
Granted Patent B2
US 11,711,255 · App. 17/463,989 · Granted Jul 25, 2023

Systems and methods for data linkage and entity resolution of continuous and un-synchronized data streams

Inventors: Pakshal Kumar H Dhelaria (Davangere, IN); Ambarish Kumar (Bangalore, IN); Saifulla Shaik (Nellore, IN); Aikaterini Kalou (Patras, GR)
Assignee: Citrix Systems, Inc.
H04L41/0233H04L41/14H04L67/562
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,711,255
App. No.
17/463,989
Granted
Jul 25, 2023
Kind
B2
Abstract

The present disclosure is directed to a scalable, extensible, fault-tolerant system for stateful joining of two or more streams that are not fully synchronized, event ordering is not guaranteed, and certain events arrive a bit late. The system can ensure to combine the events or link the data in near real-time with low latency to mitigate impacts on downstream applications, such as ML models for determining suspicious behavior. Apart from combining events, the system can ensure to propagate the needed entities to other product streams or help in entity resolution. If any of the needed data is yet to arrive, a user can configure a few parameters to achieve desired eventual and attribute consistency. The architecture is designed to be agnostic of stream processing framework and can work well with both streaming and batch paths.

Claims (58)

1. A method comprising:

creating, by one or more processors, responsive to processing a first event from a first stream of a plurality of streams, an object including first data from the first event for merging the first data with second data from a second event of a second stream of the plurality of streams;

routing, by the one or more processors, the first event to a retry stream to reprocess the first event responsive to determining that the object does not include the second data from the second event and a number of times the first event has been routed to the retry stream does not satisfy a threshold; and

transmitting, by the one or more processors, i) the first data and the second data included in the object responsive to processing the second event to update the object to include the second data or ii) the first data included in the object responsive to determining that the object does not include the second data from the second event and the number of times the first event has been routed to the retry stream satisfies the threshold.

2. The method of claim 1 , further comprising subscribing, by the one or more processors, to the first stream, the second stream, and the retry stream.

3. The method of claim 1 , further comprising:

determining, by the one or more processors, that the object does not include the second data and the number of times the first event has been routed to the retry stream does not satisfy the threshold;

updating, by the one or more processors, a count indicating the number of times the first event has been routed to the retry stream; and

wherein routing the first event comprises passing, by the one or more processors, the first event to a message broker corresponding to the retry stream.

4. The method of claim 1 , wherein creating the object comprises:

identifying, by the one or more processors, one or more attributes from the first event; and

storing, by the one or more processors, as the first data, the one or more attributes in the object.

5. The method of claim 4 , wherein creating the object for the first event comprises creating the object for the first event responsive to determining that an identifier of the first event does not match an identifier of any object stored by the one or more processors.

6. The method of claim 1 , further comprising:

receiving, by the one or more processors, the second event from the second stream;

determining, by the one or more processors, an identifier from the second event;

determining, by the one or more processors, that the identifier from the second event matches an identifier of the object; and

storing, by the one or more processors, as the second data, one or more attributes from the second event in the object.

7. The method of claim 6 , further comprising dropping, by the one or more processors, the first event responsive to transmitting either the first data stored in the object or the first data and the second data stored in the object.

8. The method of claim 1 , further comprising assigning a flag to the object indicating that the data is to be transmitted responsive to i) determining that the object does not include the second data and the number of times the first event has been routed to the retry stream satisfies the threshold or ii) determining that the object includes the second data.

9. The method of claim 1 , further comprising:

generating, by the one or more processors, responsive to determining that the object includes the second data, a third event comprising one or more attributes from the first event and the second event, and wherein transmitting the data comprises transmitting the third event; or

generating, by the one or more processors, responsive to determining that the object does not include the second data and the number of times the first event was routed to the retry stream satisfies the threshold, a fourth event comprising one or more attributes from the first event; and

wherein transmitting the data included in the object comprises transmitting the third event or the fourth event.

10. The method of claim 1 , further comprising:

updating, by the one or more processors, responsive to creating the object including the first data of the first event, the object with a first flag indicating that the first event from the first stream has been received; and

updating, by the one or more processors, responsive to updating the object to include the second data of the second event, the object with a second flag indicating that the second event from the second stream has been received.

11. A system comprising:

one or more processors configured to:

create, responsive to receiving a first event from a first stream of a plurality of streams, an object including first data from the first event for merging the first data with second data from a second event of a second stream of the plurality of streams;

route the first event to a retry stream to reprocess the first event responsive to determining that the object does not include the second data from the second event and a number of times the first event has been routed to the retry stream does not satisfy a threshold; and

transmit i) the first data and the second data included in the object responsive to processing the second event to update the object to include the second data or ii) the first data included in the object responsive to determining that the object does not include the second data from the second event and the number of times the first event has been routed to the retry stream satisfies the threshold.

12. The system of claim 11 , wherein the one or more processors are further configured to:

determine that the object does not include the second data and the number of times the first event has been routed to the retry stream does not satisfy the threshold;

update a count indicating the number of times the first event has been routed to the retry stream; and

wherein to route the first event, the one or more processors are configured to pass the first event to a message broker corresponding to the retry stream.

13. The system of claim 11 , wherein to create the object, the one or more processors are further configured to:

identify one or more attributes from the first event; and

store, as the first data, the one or more attributes in the object.

14. The system of claim 11 , wherein to create the object for the first event, the one or more processors are configured to create the object for the first event responsive to determining that an identifier of the first event does not match an identifier of any object stored by the one or more processors.

15. The system of claim 11 , wherein the one or more processors are further configured to:

receive the second event from the second stream;

determine an identifier from the second event;

determine that the identifier from the second event matches an identifier of the object; and

store, as the second data, one or more attributes from the second event in the object.

16. The system of claim 11 , wherein the one or more processors are further configured to drop the first event responsive to transmitting either the first data stored in the object or the first data and the second data stored in the object.

17. The system of claim 11 , wherein the one or more processors are further configured to assign a flag to the object indicating that the data included in the object is to be transmitted responsive to i) determining that the object does not include the second data and the number of times the first event has been routed to the retry stream satisfies the threshold or ii) determining that the object includes the second data.

18. The system of claim 11 , wherein the one or more processors are further configured to:

generate, responsive to determining that the object includes the second data, a third event comprising one or more attributes from the first event and the second event, and wherein transmitting the data comprises transmitting the third event; or

generate, responsive to determining that the object does not include the second data and the number of times the first event was routed to the retry stream satisfies the threshold, a fourth event comprising one or more attributes from the first event; and

wherein to transmit the data included in the object, the one or more processors are further configured to transmit the third event or the fourth event.

19. The system of claim 11 , wherein the one or more processors are further configured to:

update, responsive to creating the object using attributes of the first event, the object with a first flag indicating that the first event from the first stream has been received; and

update, responsive to updating the object using attributes of the second event, the object with a second flag indicating that the second event from the retry stream has been received.

20. A non-transitory computer-readable medium storing instructions that, when executed by one or more processors, cause the one or more processors to:

create, responsive to receiving a first event from a first stream of a plurality of streams, an object including first data from the first event for merging the first data with second data from a second event of a second stream of the plurality of streams;

route the first event to a retry stream to reprocess the first event responsive to determining that the object does not include the second data from the second event and a number of times the first event has been routed to the retry stream does not satisfy a threshold; and

transmit i) the first data and the second data included in the object responsive to processing the second event to update the object to include the second data or ii) the first data included in the object responsive to determining that the object does not include the second data from the second event and the number of times the first event has been routed to the retry stream satisfies the threshold.

Assignments (9)
PATENT SECURITY AGREEMENT Recorded Aug 15, 2025
From: CLOUD SOFTWARE GROUP, INC.; CITRIX SYSTEMS, INC.
To: WILMINGTON TRUST, NATIONAL ASSOCIATION, AS NOTES COLLATERAL AGENT
Reel/Frame 072488/0172 →
SECURITY INTEREST Recorded May 24, 2024
From: CLOUD SOFTWARE GROUP, INC. (F/K/A TIBCO SOFTWARE INC.); CITRIX SYSTEMS, INC.
To: WILMINGTON TRUST, NATIONAL ASSOCIATION, AS NOTES COLLATERAL AGENT
Reel/Frame 067662/0568 →
PATENT SECURITY AGREEMENT Recorded Apr 14, 2023
From: CLOUD SOFTWARE GROUP, INC. (F/K/A TIBCO SOFTWARE INC.); CITRIX SYSTEMS, INC.
To: WILMINGTON TRUST, NATIONAL ASSOCIATION, AS NOTES COLLATERAL AGENT
Reel/Frame 063340/0164 →
RELEASE AND REASSIGNMENT OF SECURITY INTEREST IN PATENT (REEL/FRAME 062113/0001) Recorded Apr 14, 2023
From: GOLDMAN SACHS BANK USA, AS COLLATERAL AGENT
To: CITRIX SYSTEMS, INC.; CLOUD SOFTWARE GROUP, INC. (F/K/A TIBCO SOFTWARE INC.)
Reel/Frame 063339/0525 →
PATENT SECURITY AGREEMENT Recorded Oct 7, 2022
From: TIBCO SOFTWARE INC.; CITRIX SYSTEMS, INC.
To: BANK OF AMERICA, N.A., AS COLLATERAL AGENT
Reel/Frame 062112/0262 →
PATENT SECURITY AGREEMENT Recorded Oct 7, 2022
From: TIBCO SOFTWARE INC.; CITRIX SYSTEMS, INC.
To: WILMINGTON TRUST, NATIONAL ASSOCIATION, AS NOTES COLLATERAL AGENT
Reel/Frame 062113/0470 →
SECOND LIEN PATENT SECURITY AGREEMENT Recorded Oct 7, 2022
From: TIBCO SOFTWARE INC.; CITRIX SYSTEMS, INC.
To: GOLDMAN SACHS BANK USA, AS COLLATERAL AGENT
Reel/Frame 062113/0001 →
SECURITY INTEREST Recorded Sep 30, 2022
From: CITRIX SYSTEMS, INC.
To: WILMINGTON TRUST, NATIONAL ASSOCIATION
Reel/Frame 062079/0001 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Sep 1, 2021
From: DHELARIA, PAKSHAL KUMAR H; KUMAR, AMBARISH; SHAIK, SAIFULLA; KALOU, AIKATERINI
To: CITRIX SYSTEMS, INC.
Reel/Frame 057359/0116 →
Continuity (2)
Continuation PCTGR2021000055 · Aug 17, 2021
Related Publication 20230055677A1 · Feb 23, 2023
Cited By (1)
US 12,335,120