IP Library Granted Patent US 12,314,264
Granted Patent B2
US 12,314,264 · App. 17/402,388 · Granted May 27, 2025

System for embedding stream processing execution in a database

Inventors: Radu Tudoran (Munich, DE); Alexander Nozdrin (Munich, DE); Stefano Bortoli (Munich, DE); Mohamad Al Hajj Hassan (Munich, DE); Cristian Axenie (Munich, DE); Hailin Li (Shenzhen, CN); Goetz Brasche (Munich, DE)
Assignee: Huawei Cloud Computing Technologies Co., Ltd.
G06F16/24568G06F16/244G06F16/2456G06F16/2457
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,314,264
App. No.
17/402,388
Granted
May 27, 2025
Kind
B2
Abstract

A database management system, comprising: a storage adapted for storing: a plurality of data objects organized according to a data model, and a plurality of stream operator wrappers each wrapping a stream operator and having at least one port for receiving, via a network, instructions for: scheduling activation of the wrapped stream operator, connecting the wrapped stream operator with another stream operator wrapped by another of the plurality of stream operator wrappers, and/or deploying the wrapped stream operator; a processing circuitry for executing the plurality of stream operator wrappers.

Claims (28)

1. A database management system, comprising:

a memory having the following stored thereon:

a plurality of data objects organized according to a data model, and

a plurality of stream operator wrappers, wherein each of the plurality of stream operator wrappers wraps a stream operator and has at least one port, wherein the at least one port is configured to receive, via a network, instructions from a remote stream engine, wherein the instructions are determined by the remote stream engine based on a plurality of messages from the plurality of stream operator wrappers, and wherein the instructions comprise at least one of scheduling activation of the wrapped stream operator, connecting the wrapped stream operator with another stream operator wrapped by another of the plurality of stream operator wrappers, and deploying the wrapped stream operator, and wherein the at least one port comprises at least one output port for outputting data to the remote stream engine, wherein the remote stream engine schedules operation of a plurality of stream operators wrapped in the plurality of stream operator wrappers; and

a processor configured to execute the plurality of stream operator wrappers, wherein the processor is configured to execute at least two of the plurality of stream operator wrappers simultaneously for processing a data flow of a common thread.

2. The database management system of claim 1 , wherein each of the plurality of stream operator wrappers includes an entry point for receiving a data flow from another of the plurality of stream operator wrappers which is interconnected thereto.

3. The database management system of claim 1 , wherein the stream operator is a member of a group consisting of the following operators: map, flatmap, reduce, join, window, filter, processfunction, groupby, keyby, shuffle, group, iterate, match, and aggregate.

4. The database management system of claim 1 , wherein the database management system is part of a distributed stream processing pipeline.

5. The database management system of claim 1 , wherein the processor executes each of at least two of the plurality of stream operator wrappers for combining a logic defined by the respective stream operator and at least two of the plurality of data objects.

6. The database management system of claim 1 , wherein the instructions are received for implementing a stream processing pipeline integration.

7. The database management system of claim 1 , wherein the instructions are received to control a job execution for enabling coordination and distribution of jobs according to a stream processing plan.

8. The database management system of claim 7 , wherein the stream processing plan causes the processor to build a stream application topology.

9. The database management system of claim 1 , wherein the stream operator is a single thread application defining logic for processing data from the plurality of data objects.

10. The database management system of claim 1 , wherein the stream operator is a user defined function.

11. The database management system of claim 1 , wherein the at least one port comprises at least one output port for outputting data to the remote stream engine, wherein the remote stream engine supervises the execution of the respective stream operator.

12. A stream engine, comprising:

an interface in communication, via a network, with a remote database management system, wherein the interface is configured to receive a plurality of messages from a plurality of stream operator wrappers executed on the remote database management system, wherein the remote database management system comprises a memory having a plurality of data objects and the plurality of stream operator wrappers stored thereon, and wherein each of the plurality of messages is indicative of an outcome of executing a stream operator wrapped by one of the plurality of stream operator wrappers, and wherein the stream engine schedules operation of a plurality of stream operators wrapped in the plurality of stream operator wrappers; and

a processor configured to execute a policy for scheduling or interconnecting at least some of a plurality of stream operators wrapped in the plurality of stream operator wrappers according to data in the plurality of messages, and to send instructions to the remote database management system based on the executed policy, wherein the processor is configured to execute at least two of the plurality of stream operator wrappers simultaneously for processing a data flow of a common thread.

13. The stream engine of claim 12 , wherein the processor is configured to instruct a deployment of one or more of the plurality of stream operators by sending an indication to one or more of the plurality of stream operator wrappers via the network.

14. A system, comprising:

a stream engine; and

a database management sub-system comprising a processor and a memory having the following stored thereon:

a plurality of data objects organized according to a data model; and

a plurality of stream operator wrappers, wherein each of the plurality of stream operator wrappers wraps a stream operator and has at least one port;

wherein at least two of the plurality of stream operator wrappers are executed by the processor based on instructions, wherein the instructions are received through the at least one port from the stream engine via a network, wherein the instructions are determined by the stream engine based on a plurality of messages from the plurality of stream operator wrappers, and wherein the instructions comprise at least one of scheduling activation of the wrapped stream operator, connecting the wrapped stream operator with another stream operator wrapped by another of the plurality of stream operator wrappers, and deploying the wrapped stream operator, and wherein the at least one port comprises at least one output port for outputting data to the stream engine, wherein the stream engine schedules operation of a plurality of stream operators wrapped in the plurality of stream operator wrappers, wherein the processor is configured to execute at least two of the plurality of stream operator wrappers simultaneously for processing a data flow of a common thread.

15. The system of claim 14 , wherein each of the plurality of stream operator wrappers includes an entry point for receiving a data flow from another of the plurality of stream operator wrappers which is interconnected thereto.

16. The system of claim 14 , wherein the stream operator is a member of a group consisting of the following operators: map, flatmap, reduce, join, window, filter, processfunction, groupby, keyby, shuffle, group, iterate, match, and aggregate.

17. The system of claim 14 , wherein the database management sub-system is part of a distributed stream processing pipeline.

Assignments (2)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Apr 30, 2025
From: TUDORAN, RADU; NOZDRIN, ALEXANDER; BORTOLI, STEFANO; HASSAN, MOHAMAD AL HAJJ; AXENIE, CRISTIAN; LI, HAILIN; BRASCHE, GOETZ
To: HUAWEI TECHNOLOGIES CO., LTD.
Reel/Frame 070990/0165 →
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Mar 1, 2022
From: HUAWEI TECHNOLOGIES CO., LTD.
To: HUAWEI CLOUD COMPUTING TECHNOLOGIES CO., LTD.
Reel/Frame 059267/0088 →
Continuity (2)
Continuation PCTEP2019053795 · Feb 15, 2019
Related Publication 20210374144A1 · Dec 2, 2021
References Cited (35)
US 7849507B1 · Bloch · 2010 [cited by examiner]
US 8805819B1 · Black · 2014 [cited by applicant]
US 9183339B1 · Shirazi · 2015 [cited by examiner]
US 9607045B2 · Fisher et al. · 2017 [cited by applicant]
US 10180861B2 · Raghavan · 2019 [cited by examiner]
US 10389764B2 · Johnson · 2019 [cited by examiner]
US 20010056492A1 · Bressoud · 2001 [cited by examiner]
US 20020133504A1 · Vlahos · 2002 [cited by examiner]
US 20040243550A1 · Gu · 2004 [cited by examiner]
US 20070067274A1 · Han · 2007 [cited by examiner]
US 20100030896A1 · Chandramouli · 2010 [cited by examiner]
US 20120078868A1 · Chen · 2012 [cited by examiner]
US 20120324453A1 · Chandramouli · 2012 [cited by examiner]
US 20140032525A1 · Merriman · 2014 [cited by examiner]
US 20150248461A1 · Theeten · 2015 [cited by examiner]
US 20150317364A1 · Branson et al. · 2015 [cited by applicant]
US 20150372882A1 · Qian · 2015 [cited by examiner]
US 20160180022A1 · Paixao · 2016 [cited by examiner]
US 20170060947A1 · Zhang · 2017 [cited by examiner]
US 20170124151A1 · Ji · 2017 [cited by examiner]
US 20180227280A1 · Barsness et al. · 2018 [cited by applicant]
US 20180300381A1 · Horowitz · 2018 [cited by examiner]
US 20190042290A1 · Bailey · 2019 [cited by examiner]
US 20190124007A1 · Cook · 2019 [cited by examiner]
US 20190132387A1 · Singh · 2019 [cited by examiner]
US 20190207990A1 · Cook · 2019 [cited by examiner]
US 20190251081A1 · Branson · 2019 [cited by examiner]
US 20190377817A1 · McCluskey · 2019 [cited by examiner]
US 20200092314A1 · Mital · 2020 [cited by examiner]
US 20200117757A1 · Yanamandra · 2020 [cited by examiner]
CN 105308592A · 2016 [cited by applicant]
CN 109074377A · 2018 [cited by applicant]
Ji et al., “Optimization of Continuous Queries in Federated Database and Stream Processing Systems,” https://www.slideshare.net/zbigniew.jerzak/2015-0305-btwtalkv2, pp. 1-14 (Mar. 16, 2015). [cited by applicant]
“How Continuous Querying Works,” Apache Geode, https://geode.apache.org/docs/guide/114/developing/continuous_querying/how_continuous_querying_works.html, Total 2 pages (2021). [cited by applicant]
Ji et al., “Optimization of Continuous Queries in Federated Database and Stream Processing Systems,” BTW, Computer Science, pp. 403-422 (2015). [cited by applicant]