# 𝐖𝐡𝐚𝐭 𝐢𝐬 𝐒𝐭𝐫𝐮𝐜𝐭𝐮𝐫𝐞𝐝 𝐒𝐭𝐫𝐞𝐚𝐦𝐢𝐧𝐠?

𝐖𝐡𝐚𝐭 𝐢𝐬 𝐒𝐭𝐫𝐮𝐜𝐭𝐮𝐫𝐞𝐝 𝐒𝐭𝐫𝐞𝐚𝐦𝐢𝐧𝐠?

Think of Structured Streaming as a way to handle data that's constantly arriving, like a river of information. Spark processes this data in small chunks, or "micro-batches," making it feel like it's happening instantly. It's built on top of the Spark SQL engine, so you can use familiar SQL queries to process your data.

𝐊𝐞𝐲 𝐂𝐨𝐧𝐜𝐞𝐩𝐭𝐬 𝐄𝐱𝐩𝐥𝐚𝐢𝐧𝐞𝐝 𝐬𝐢𝐦𝐩𝐥𝐲:

𝐈𝐧𝐩𝐮𝐭 𝐒𝐨𝐮𝐫𝐜𝐞𝐬: This is where your data comes from. Common sources include Kafka (a distributed messaging system), file directories (like S3 or HDFS), and network sockets.

𝙍𝙚𝙖𝙡-𝙬𝙤𝙧𝙡𝙙 𝙀𝙭𝙖𝙢𝙥𝙡𝙚: Reading user activity logs from Kafka topics to track website clicks in real-time.

Code Sample (Python): spark.readStream.format("kafka").option("kafka.bootstrap.servers", "host:port").option("subscribe", "topic\_name").load()

𝐓𝐫𝐚𝐧𝐬𝐟𝐨𝐫𝐦𝐚𝐭𝐢𝐨𝐧𝐬 & 𝐀𝐜𝐭𝐢𝐨𝐧𝐬: This is how you process and manipulate your streaming data.

𝙏𝙧𝙖𝙣𝙨𝙛𝙤𝙧𝙢𝙖𝙩𝙞𝙤𝙣𝙨: Define what operations you want to perform, like filtering, joining, or aggregating data. They are "lazy" and don't execute until you trigger an action.

Real-world Example: Filtering out irrelevant logs or joining user activity data with user profile information.

Code Sample (Python): df.filter("user\_id > 100").groupBy("action").count()

𝘼𝙘𝙩𝙞𝙤𝙣𝙨: Trigger the processing of data and define where the output goes. Actions are essential for writing the results to a sink.

Real-world Example: Writing the processed data to a database or displaying results on a dashboard.

Code Sample (Python): query = df.writeStream.outputMode("append").format("console").start()

𝐒𝐢𝐧𝐤𝐬: This is where you want to store or display your processed data. Options include console, Delta Lake, or various databases.

Real-world Example: Writing processed analytics data to Delta Lake for further analysis or to a console for debugging.

Code Sample (Python): query = df.writeStream.outputMode("append").format("delta").option("path", "/path/to/delta/table").start()

𝐎𝐮𝐭𝐩𝐮𝐭 𝐌𝐨𝐝𝐞𝐬: This determines how Spark writes the results to the sink.

𝘼𝙥𝙥𝙚𝙣𝙙 𝙈𝙤𝙙𝙚: Only new rows added to the result table are written.

𝘾𝙤𝙢𝙥𝙡𝙚𝙩𝙚 𝙈𝙤𝙙𝙚: The entire result table is written every time.

𝙐𝙥𝙙𝙖𝙩𝙚 𝙈𝙤𝙙𝙚: Only the rows that have been updated since the last trigger are written.

Real-world Example: Append for new user signups, Update for profile updates, Complete for a top 10 list.

Code Sample (Python): query = df.writeStream.outputMode("append")...

𝐓𝐫𝐢𝐠𝐠𝐞𝐫𝐬: These control when Spark processes a new batch of data. Options include fixed intervals or one-time execution.

Real-world Example: Processing data every minute or once an hour.

Code Sample (Python): query = df.writeStream.trigger(processingTime="1 minute")...

𝐄𝐯𝐞𝐧𝐭 𝐓𝐢𝐦𝐞 𝐏𝐫𝐨𝐜𝐞𝐬𝐬𝐢𝐧𝐠: Handling data based on when it was generated, not when Spark received it. This is crucial for accurate analysis, especially when dealing with late-arriving data.

Real-world Example: Aggregating user actions by their actual occurrence time to understand user behavior accurately.

Code Sample (Python): withWatermark("event\_time", "5 minutes").groupBy(window("event\_time", "1 minute")).count()
