IP Library Granted Patent US 11,757,959
Granted Patent B2
US 11,757,959 · App. 17/167,470 · Granted Sep 12, 2023

Dynamic data stream processing for Apache Kafka using GraphQL

Inventors: Wojciech Trocki (Waterford, IE); Manyanda Chitimbo (Creteil, FR); Enda Martin Phelan (Waterford, IE)
Assignee: Red Hat, Inc.
H04L65/61G06F16/245
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,757,959
App. No.
17/167,470
Granted
Sep 12, 2023
Kind
B2
Abstract

Systems and methods for specifying a stream processing topology via a client-side API without server-side support. A schema may be generated by a client-side application using a query language and transmitted to a stream processor registry, wherein the schema defines a desired data stream. The stream processor registry, acting as a server-side run time corresponding to the query language, may store the schema as metadata. The stream processor registry may generate a stream processing topology based on the metadata to obtain data relevant to the data stream and generate a user-specific topic comprising the data relevant to the data stream. The stream processor registry may filter the data relevant to the data stream based a subscription call by the client to generate a target topic comprising portions of the data relevant to the data stream.

Claims (43)

1. A method comprising:

receiving a schema defined by a client, the schema comprising one or more mutations that collectively define a data stream, wherein each of the one or more mutations comprise a set of rules that define data the client wants to search for;

reading one or more topics provided by a data streaming platform to obtain data relevant to the data stream based on the one or more mutations;

generating, by a processing device, a single user-specific topic comprising the data relevant to the data stream; and

providing the client access to the user-specific topic.

2. The method of claim 1 , wherein providing the client access to the user-specific topic comprises:

transmitting portions of the data relevant to the data stream to the client based on a call to a subscription to the user-specific topic by the client.

3. The method of claim 2 , further comprising:

filtering the data relevant to the data stream based on one or more of: offset rules, filtering rules, aggregation rules, and windowing rules specified by the call to generate a target topic comprising the portions of the data relevant to the data stream.

4. The method of claim 1 , wherein the schema is received by a stream processor registry that functions as a server-side run time to obtain the data relevant to the data stream by executing functions specified by the one or more mutations.

5. The method of claim 4 , further comprising:

storing the one or more mutations as metadata with the stream processor registry, wherein the stream processor registry utilizes the metadata to build a stream processor to obtain the data relevant to the data stream in response to a call to a subscription to the user-specific topic by the client.

6. The method of claim 4 , wherein the schema is defined using a query language corresponding to the server-side run time.

7. The method of claim 6 , wherein the query language is GraphQL.

8. A system comprising:

a memory; and

a processing device operatively coupled to the memory, the processing device to:

receive a schema defined by a client, the schema comprising one or more mutations that collectively define a data stream, wherein each of the one or more mutations comprise a set of rules that define data the client wants to search for;

read one or more topics provided by a data streaming platform to obtain data relevant to the data stream based on the one or more mutations;

generate a single user-specific topic comprising the data relevant to the data stream; and

provide the client access to the user-specific topic.

9. The system of claim 8 , wherein to provide the client access to the user-specific topic, the processing device is to:

transmit portions of the data relevant to the data stream to the client based on a call to a subscription to the user-specific topic by the client.

10. The system of claim 9 , wherein the processing device is further to:

filter the data relevant to the data stream based on one or more of: offset rules, filtering rules, aggregation rules, and windowing rules specified by the call to generate a target topic comprising the portions of the data relevant to the data stream.

11. The system of claim 8 , wherein the schema is received by a stream processor registry that functions as a server-side run time to obtain the data relevant to the data stream by executing functions specified by the one or more mutations.

12. The system of claim 11 , wherein the processing device is further to:

store the one or more mutations as metadata with the stream processor registry, wherein the stream processor registry utilizes the metadata to build a stream processor to obtain the data relevant to the data stream in response to a call to a subscription to the user-specific topic by the client.

13. The system of claim 11 , wherein the schema is defined using a query language corresponding to the server-side run time.

14. The system of claim 13 , wherein the query language is GraphQL.

15. A non-transitory computer-readable medium, having instructions stored thereon which, when executed by a processing device, cause the processing device to:

receive a schema defined by a client, the schema comprising one or more mutations that collectively define a data stream, wherein each of the one or more mutations comprise a set of rules that define data the client wants to search for;

read one or more topics provided by a data streaming platform to obtain data relevant to the data stream based on the one or more mutations;

generate, by the processing device, a single user-specific topic comprising the data relevant to the data stream; and

provide the client access to the user-specific topic.

16. The non-transitory computer-readable medium of claim 15 , wherein to provide the client access to the user-specific topic, the processing device is to:

transmit portions of the data relevant to the data stream to the client based on a call to a subscription to the user-specific topic by the client.

17. The non-transitory computer-readable medium of claim 16 , wherein the processing device is further to:

filter the data relevant to the data stream based on one or more of: offset rules, filtering rules, aggregation rules, and windowing rules specified by the call to generate a target topic comprising the portions of the data relevant to the data stream.

18. The non-transitory computer-readable medium of claim 15 , wherein the schema is received by a stream processor registry that functions as a server-side run time to obtain the data relevant to the data stream by executing functions specified by the one or more mutations.

19. The non-transitory computer-readable medium of claim 18 , wherein the processing device is further to:

store the one or more mutations as metadata with the stream processor registry, wherein the stream processor registry utilizes the metadata to build a stream processor to obtain the data relevant to the data stream in response to a call to a subscription to the user-specific topic by the client.

20. The non-transitory computer-readable medium of claim 18 , wherein the schema is defined using a query language corresponding to the server-side run time.

Assignments (2)
CHANGE OF NAME Recorded Mar 3, 2026
From: RED HAT, INC.
To: RED HAT, LLC
Reel/Frame 074913/0759 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Feb 4, 2021
From: TROCKI, WOJCIECH; CHITIMBO, MANYANDA; PHELAN, ENDA MARTIN
To: RED HAT, INC.
Reel/Frame 055149/0158 →
Cited By (1)
US 12,326,776