IP Library Granted Patent US 11,755,618
Granted Patent B2
US 11,755,618 · App. 15/913,915 · Granted Sep 12, 2023

Stateless stream handling and resharding

Inventors: Ori Modai (Ramat Hasharon, IL); Orit Nissan-Messing (Hod Hasharon, IL); Yaron Haviv (Tel Mond, IL)
Assignee: Iguazio Systems Ltd.
G06F16/278G06F16/2365G06F16/24568
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,755,618
App. No.
15/913,915
Granted
Sep 12, 2023
Kind
B2
Abstract

Systems and methods are disclosed for stateless stream handling and resharding. In one implementation, a first shard comprising one or more messages is generated. The first shard is associated with a first state attribute. The first shard and the first state attribute are provided as an update within a data stream. In another implementation, a first shard including a first state attribute is received within a first stream. A message that is inconsistent with the first state attribute is identified within the first shard. The message is associated as an attribute of the first shard. A second shard including a second state attribute is received. Based on the second state attribute, a position of the message within the second shard is determined. The message is inserted into the second shard based on the determining.

Claims (43)

1. A system comprising:

a processing device; and

a memory coupled to the processing device and storing instructions that, when executed by the processing device, cause the system to perform operations comprising:

generating a first shard comprising one or more messages;

associating the first shard with a first state attribute;

providing the first shard and the first state attribute as an atomic update within a data stream;

requesting the first state attribute from the first shard, wherein the first state attribute comprises a token that reflects a processing capacity of a streaming system;

receiving the first state attribute;

providing a second shard within the data stream based on the received first state attribute; and

adjusting an operation of a message production source of at least one of the one or more messages within the streaming system based on the received token that reflects the processing capacity of the streaming system.

2. The system of claim 1 , wherein the first state attribute reflects an importance of one or more of the messages.

3. The system of claim 1 , wherein the first state attribute reflects a location of one or more of the messages.

4. The system of claim 1 , wherein an operation associated with the first shard is performed based on the first state attribute.

5. The system of claim 1 , wherein the atomic update comprises a plurality of updates that are collectively performed or rejected.

6. The system of claim 1 , wherein providing the first shard and the first state attribute comprises providing the first shard and the first state attribute as a conditional update within the data stream.

7. The system of claim 1 , wherein the first state attribute comprises a first sequence identifier.

8. A non-transitory computer readable medium having instructions stored thereon that, when executed by a processing device, cause the processing device to perform operations comprising:

generating a first shard comprising one or more messages;

associating the first shard with a first state attribute;

providing the first shard and the first state attribute as an update within a data stream;

requesting the first state attribute from the first shard, wherein the first state attribute comprises a token that reflects a processing capacity of a streaming system;

receiving the first state attribute;

providing a second shard within the data stream based on the received first state attribute; and

adjusting an operation of a message production source of at least one of the one or more messages within the streaming system based on the received token that reflects the processing capacity of the streaming system.

9. The non-transitory computer-readable medium of claim 8 , wherein an operation associated with the first shard is performed based on the first state attribute.

10. The non-transitory computer-readable medium of claim 8 , wherein providing the first shard and the first state attribute comprises providing the first shard and the first state attribute as at least one of (a) an atomic update or (b) a conditional update within the data stream.

11. A method comprising:

generating a first shard comprising one or more messages;

associating the first shard with a first state attribute;

providing the first shard and the first state attribute as an atomic update within a data stream;

requesting the first state attribute from the first shard, wherein the first state attribute comprises a token that reflects a processing capacity of a streaming system;

receiving the first state attribute;

providing a second shard within the data stream based on the received first state attribute; and

adjusting an operation of a message production source of at least one of the one or more messages within the streaming system based on the received token that reflects the processing capacity of the streaming system.

12. The method of claim 11 , wherein the first state attribute reflects an importance of one or more of the messages.

13. The method of claim 11 , wherein the first state attribute reflects a location of one or more of the messages.

14. The method of claim 11 , wherein an operation associated with the first shard is performed based on the first state attribute.

15. The method of claim 11 , wherein the atomic update comprises a plurality of updates that are collectively performed or rejected.

16. The non-transitory computer-readable medium of claim 8 , wherein the first state attribute reflects an importance of one or more of the messages.

17. The non-transitory computer-readable medium of claim 8 , wherein the first state attribute reflects a location of one or more of the messages.

18. The non-transitory computer-readable medium of claim 8 , wherein the update comprises a plurality of updates that are collectively performed.

19. The non-transitory computer-readable medium of claim 8 , wherein the update comprises a plurality of updates that are collectively rejected.

20. The non-transitory computer-readable medium of claim 8 , wherein the first state attribute comprises a first sequence identifier.

Assignments (2)
SECURITY INTEREST Recorded May 6, 2019
From: IGUAZIO SYSTEMS LTD.
To: KREOS CAPITAL VI (EXPERT FUND) L.P.
Reel/Frame 049086/0455 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 11, 2018
From: MODAI, ORI; NISSAN-MESSING, ORIT; HAVIV, YARON
To: IGUAZIO SYSTEMS LTD.
Reel/Frame 045170/0531 →
Continuity (1)
Related Publication 20190278863A1 · Sep 12, 2019