350
N. Balac
Fig. 6.18 Spark Streaming processing
Fig. 6.19 Spark Streaming functionality
memory. The Spark engine processes the batches, while optimizing latency, and
outputs the results to external systems as shown in Fig. 6.19.
Spark Streaming maintains a state based on data coming in a stream often
referred to as stateful computations. In addition, Spark Streaming allows window
operations where a specified time frame could be used to perform operations on
the data. The sliding time interval in the window is used for updating the window,
utilizing the window length and sliding interval parameters. When the window slides
over a source DStream, the underlying RDDs are combined and operated upon to
produce the RDDs of the windowed DStream [28]. Spark tasks are assigned to the
workers dynamically on the basis of data locality and available resources, therefore
optimizing load balancing and fault recovery.
Spark Streaming’s data stream can originate from the source data stream or the
processed data stream generated by transforming the input stream. Internally, a
DStream is represented by a continuous series of RDDs. Every input DStream is
associated with a Receiver, which receives the data from a source and stores it in
executor memory.
Analogous to Spark RDDs, Spark transformations enable DStream modifications. Input DStreams support many transformations that are applicable to RDDs,
including map, filter, count, countbyvalue, reduce, union, etc. Spark Streaming
enables two categories of built-in streaming sources: basic and advanced sources.
Basic sources are typically directly available in the StreamingContext API, like file
systems, and socket connections. Advanced sources typically include Kafka, Flume,
N. Balac
Fig. 6.18 Spark Streaming processing
Fig. 6.19 Spark Streaming functionality
memory. The Spark engine processes the batches, while optimizing latency, and
outputs the results to external systems as shown in Fig. 6.19.
Spark Streaming maintains a state based on data coming in a stream often
referred to as stateful computations. In addition, Spark Streaming allows window
operations where a specified time frame could be used to perform operations on
the data. The sliding time interval in the window is used for updating the window,
utilizing the window length and sliding interval parameters. When the window slides
over a source DStream, the underlying RDDs are combined and operated upon to
produce the RDDs of the windowed DStream [28]. Spark tasks are assigned to the
workers dynamically on the basis of data locality and available resources, therefore
optimizing load balancing and fault recovery.
Spark Streaming’s data stream can originate from the source data stream or the
processed data stream generated by transforming the input stream. Internally, a
DStream is represented by a continuous series of RDDs. Every input DStream is
associated with a Receiver, which receives the data from a source and stores it in
executor memory.
Analogous to Spark RDDs, Spark transformations enable DStream modifications. Input DStreams support many transformations that are applicable to RDDs,
including map, filter, count, countbyvalue, reduce, union, etc. Spark Streaming
enables two categories of built-in streaming sources: basic and advanced sources.
Basic sources are typically directly available in the StreamingContext API, like file
systems, and socket connections. Advanced sources typically include Kafka, Flume,
