Dynamic reduction of stream backpressure
First Claim
1. A computer program product for processing data, the computer program product comprising:
- a non-transitory computer-readable storage medium having computer-readable program code embodied therewith, the computer-readable program code comprising computer-readable program code configured to;
receive streaming data to be processed by a plurality of processing elements, wherein each of the processing elements is an executable portion of code;
establish an operator graph of the plurality of processing elements, the operator graph defining at least one execution path through which the streaming data flows through the plurality of processing elements, wherein each processing element in the execution path is configured to at least one of receive data from at least one upstream processing element and transmit data to at least one downstream processing element;
monitor an indicator of backpressure associated with a first processing element in the execution path, wherein the indicator of backpressure represents the data throughput of the first processing element;
prioritize at least two of the plurality of processing elements that are one of upstream and downstream of the first processing element in the execution path, wherein priority is based on at least one of (i) an importance of a job associated with each of the at least two of the processing elements, (ii) an importance of data transmitted by each of the at least two of the processing elements, and (iii) a transient time required by each of the processing elements in the at least two of the processing elements to process the streaming data ; and
upon determining that the indicator of backpressure satisfies a predetermined threshold indicating that the first processing element is processing data at an output rate that is less than an input rate at which the first processing element receives data, perform a corrective action to change the rate of data flowing through the first processing element to enable the first processing element to process data such that the output rate is greater than or equal to the input rate, wherein the corrective action is performed on a lowest priority processing element of the prioritized processing elements.
1 Assignment
0 Petitions
Accused Products
Abstract
Techniques are described for eliminating backpressure in a distributed system by changing the rate data flows through a processing element. Backpressure occurs when data throughput in a processing element begins to decrease, for example, if new processing elements are added to the operating chart or if the distributed system is required to process more data. Indicators of backpressure (current or future) may be monitored. Once current backpressure or potential backpressure is identified, the operator graph or data rates may be altered to alleviate the backpressure. For example, a processing element may reduce the data rates it sends to processing elements that are downstream in the operator graph, or processing elements and/or data paths may be eliminated. In one embodiment, processing elements and associate data paths may be prioritized so that more important execution paths are maintained. In another embodiment, if a request to add one or more processing elements may cause future backpressure, the request may be refused.
-
Citations
18 Claims
-
1. A computer program product for processing data, the computer program product comprising:
-
a non-transitory computer-readable storage medium having computer-readable program code embodied therewith, the computer-readable program code comprising computer-readable program code configured to; receive streaming data to be processed by a plurality of processing elements, wherein each of the processing elements is an executable portion of code; establish an operator graph of the plurality of processing elements, the operator graph defining at least one execution path through which the streaming data flows through the plurality of processing elements, wherein each processing element in the execution path is configured to at least one of receive data from at least one upstream processing element and transmit data to at least one downstream processing element; monitor an indicator of backpressure associated with a first processing element in the execution path, wherein the indicator of backpressure represents the data throughput of the first processing element; prioritize at least two of the plurality of processing elements that are one of upstream and downstream of the first processing element in the execution path, wherein priority is based on at least one of (i) an importance of a job associated with each of the at least two of the processing elements, (ii) an importance of data transmitted by each of the at least two of the processing elements, and (iii) a transient time required by each of the processing elements in the at least two of the processing elements to process the streaming data ; and upon determining that the indicator of backpressure satisfies a predetermined threshold indicating that the first processing element is processing data at an output rate that is less than an input rate at which the first processing element receives data, perform a corrective action to change the rate of data flowing through the first processing element to enable the first processing element to process data such that the output rate is greater than or equal to the input rate, wherein the corrective action is performed on a lowest priority processing element of the prioritized processing elements. - View Dependent Claims (2, 3, 8, 9)
-
-
4. A system for processing data, comprising:
-
a computer processor; and a memory containing a program that, when executed on the computer processor, performs an operation for processing data, comprising; receiving streaming data to be processed by a plurality of processing elements, wherein each of the processing elements is an executable portion of code; establishing an operator graph of the plurality of processing elements, the operator graph defining at least one execution path through which the streaming data flows through the plurality of processing elements, wherein each processing element in the execution path is configured to at least one of receive data from at least one upstream processing element and transmit data to at least one downstream processing element; monitoring an indicator of backpressure associated with a first processing element in the execution path, wherein the indicator of backpressure represents the data throughput of the first processing element; prioritizing at least two of the plurality of processing elements that are one of upstream and downstream of the first processing element in the execution path, wherein priority is based on at least one of (i) an importance of a job associated with each of the at least two of the processing elements, (ii) an importance of data transmitted by each of the at least two of the processing elements, and (iii) a transient time required by each of the processing elements in the at least two of the processing elements to process the streaming data ; and upon determining that the indicator of backpressure satisfies a predetermined threshold indicating that the first processing element is processing data at an output rate that is less than an input rate at which the first processing element receives data, performing a corrective action to change the rate of data flowing through the first processing element to enable the first processing element to process data such that the output rate is greater than or equal to the input rate, wherein the corrective action is performed on a lowest priority processing element of the prioritized processing elements. - View Dependent Claims (5, 6, 7, 10, 11)
-
-
12. A method of processing data, comprising:
-
receiving streaming data to be processed by a plurality of processing elements, the processing elements processing at least a portion of the received data by operation of one or more computer processors, wherein each of the processing elements is an executable portion of code; establishing an operator graph of the plurality of processing elements, the operator graph defining at least one execution path through which the streaming data flows through the plurality of processing elements, wherein each processing element in the execution path is configured to at least one of receive data from at least one upstream processing element and transmit data to at least one downstream processing element; monitoring an indicator of backpressure associated with a first processing element in the execution path, wherein the indicator of backpressure represents the data throughput of the first processing element; prioritizing at least two of the plurality of processing elements that are one of upstream and downstream of the first processing element in the execution path, wherein priority is based on at least one of (i) an importance of a job associated with each of the at least two of the processing elements, (ii) an importance of data transmitted by each of the at least two of the processing elements, and (iii) a transient time required by each of the processing elements in the at least two of the processing elements to process the streaming data ; and upon determining that the indicator of backpressure satisfies a predetermined threshold indicating that the first processing element is processing data at an output rate that is less than an input rate at which the first processing element receives data, perform a corrective action to change the rate of data flowing through the first processing element to enable the first processing element to process data such that the output rate is greater than or equal to the input rate, wherein the corrective action is performed on a lowest priority processing element of the prioritized processing elements. - View Dependent Claims (13, 14, 15, 16, 17, 18)
-
Specification