0% found this document useful (0 votes)
17 views12 pages

Spark Structured Streaming Overview

The document outlines the concepts and components of Structured Streaming in Apache Spark, focusing on processing streaming data through DataStreamReader and DataStreamWriter. It discusses the treatment of infinite data as a table, trigger intervals for processing data, output modes, checkpointing for fault tolerance, and guarantees for stream processing. Additionally, it highlights unsupported operations in streaming DataFrames, such as sorting and deduplication.

Uploaded by

xiyipix919
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd
0% found this document useful (0 votes)
17 views12 pages

Spark Structured Streaming Overview

The document outlines the concepts and components of Structured Streaming in Apache Spark, focusing on processing streaming data through DataStreamReader and DataStreamWriter. It discusses the treatment of infinite data as a table, trigger intervals for processing data, output modes, checkpointing for fault tolerance, and guarantees for stream processing. Additionally, it highlights unsupported operations in streaming DataFrames, such as sorting and deduplication.

Uploaded by

xiyipix919
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd

Structured Streaming

Learning Objectives

u Process streaming data

u DataStreamReader

u DataStreamWriter
Data Stream

u Any data source that grows over time

u New files landing in cloud storage

u Updates to a database captured in a CDC feed

u Events queued in a pub/sub messaging feed


Processing Data Stream

u 2 approaches:

1. Reprocess the entire source dataset each time

2. Only process those new data added since last update


u Structured Streaming
Spark Structured Streaming

infinite data source

data sink
Treating Infinite Data as a Table

Input Data Stream Unbounded Table


Input Streaming Table
Input_Table Output_Table

streamDF

streamDF = [Link] [Link]


.table("Input_Table") .trigger(processingTime="2 minutes")
.outputMode("append")
.option("checkpointLocation", "/path")
.table("Output_Table")
Trigger Intervals
[Link]
.trigger(processingTime="2 minutes")
.outputMode("append")
.option("checkpointLocation", "/path")
.table(”Output_Table")

Trigger Method call Behavior


Unspecified Default: processingTime="500ms"

Fixed interval .trigger(processingTime=”5 minutes") Process data in micro-batches at


the user-specified intervals

Triggered .trigger(once=True) Process all available data in a


batch single batch, then stop

Triggered .trigger(availableNow=True) Process all available data in


micro-batches multiple micro-batches, then stop
Output Modes
[Link]
.trigger(processingTime="2 minutes")
.outputMode("append")
.option("checkpointLocation", "/path")
.table(”Output_Table")

Mode Method call Behavior

Append .outputMode("append") Only newly appended rows are incrementally


(Default) appended to the target table with each batch

Complete .outputMode("complete") The target table is overwritten with each batch


Checkpointing
[Link]
.trigger(processingTime="2 minutes")
.outputMode("append")
.option("checkpointLocation", "/path")
.table(”Output_Table")

u Store stream state

u Track the progress of your stream processing

u Can Not be shared between separate streams


Guarantees

1. Fault Tolerance
u Checkpointing + Write-ahead logs

u record the offset range of data being processed during each trigger interval.

2. Exactly-once guarantee
u Idempotent sinks
Unsupported Operations

u Some operations are not supported by streaming DataFrame


u Sorting
u Deduplication

u Advanced methods
u Windowing
u Watermarking

You might also like