Skip to main content

Command Palette

Search for a command to run...

๐–๐ก๐š๐ญ ๐ข๐ฌ ๐’๐ญ๐ซ๐ฎ๐œ๐ญ๐ฎ๐ซ๐ž๐ ๐’๐ญ๐ซ๐ž๐š๐ฆ๐ข๐ง๐ ?

Updated
โ€ข3 min readโ€ขView as Markdown
A

Hi, I am Aatish Raj Having Extensive Experience in Bigdata

๐Ÿš€I Have good knowledge of Hadoop and it's internals.

๐Ÿš€I have good knowledge of ingestion tools like Sqoop

๐Ÿš€I have good knowledge of dataWare Houses like Hive

๐Ÿš€I have Good knowledge of๐Ÿ”ฅ Spark with Scala(Dataframes, Datasets, SparkSql) and it's internals

๐Ÿš€I have good knowlege over AWS(EMR, S3,Glue)

โœ๏ธTalks About #Data-Engineering

โœ๏ธTalks about SQL

A technology enthusiast and problem-solver, I specialize in Hadoop, MapReduce, Sqoop, Hive, Spark, AWS, SQL, Scala, Datastructures, and Algorithms. I have successfully designed and implemented solutions for diverse projects. My expertise in designing, coding, and troubleshooting allows me to quickly develop solutions and provide effective solutions to challenging problems. With a proven track record of success, I am well-equipped to take on new projects and deliver results

๐–๐ก๐š๐ญ ๐ข๐ฌ ๐’๐ญ๐ซ๐ฎ๐œ๐ญ๐ฎ๐ซ๐ž๐ ๐’๐ญ๐ซ๐ž๐š๐ฆ๐ข๐ง๐ ?

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()

More from this blog

๐’๐ญ๐จ๐ฉ ๐ญ๐ซ๐ž๐š๐ญ๐ข๐ง๐  ๐‡๐ฎ๐๐ข ๐ฏ๐ฌ. ๐ˆ๐œ๐ž๐›๐ž๐ซ๐  ๐š๐ฌ ๐š๐ง '๐ž๐ข๐ญ๐ก๐ž๐ซ/๐จ๐ซ' ๐›๐š๐ญ๐ญ๐ฅ๐ž.

In production architectures, they solve two entirely different bottlenecks: โšก ๐ˆ๐ง๐ ๐ž๐ฌ๐ญ๐ข๐จ๐ง & ๐‚๐ƒ๐‚ (๐‡๐ฎ๐๐ข) ๐ŸงŠ ๐‹๐š๐ซ๐ ๐ž-๐ฌ๐œ๐š๐ฅ๐ž ๐Ž๐‹๐€๐ & ๐๐ˆ (๐ˆ๐œ๐ž๐›๐ž๐ซ๐ ) Here is how they fit toge

Aug 25, 20262 min read

๐Ÿ๐Ÿ ๐๐š๐ญ๐š๐›๐š๐ฌ๐ž ๐ค๐ž๐ฒ๐ฌ we all need to know

โ†˜๏ธ ๐Œ๐จ๐ฌ๐ญ ๐๐ž๐ฏ๐ž๐ฅ๐จ๐ฉ๐ž๐ซ๐ฌ ๐ฆ๐ข๐ฑ ๐ญ๐ก๐ž๐ฌ๐ž ๐ฎ๐ฉ ๐๐ฎ๐ซ๐ข๐ง๐  ๐ข๐ง๐ญ๐ž๐ซ๐ฏ๐ข๐ž๐ฐ๐ฌ. ๐Ÿš€ Here is a quick cheat sheet of all ๐Ÿ๐Ÿ ๐๐š๐ญ๐š๐›๐š๐ฌ๐ž ๐ค๐ž๐ฒ๐ฌ we all need to know: Core Identification

Aug 23, 20263 min read

๐„๐ง๐-๐ญ๐จ-๐„๐ง๐ ๐ƒ๐š๐ญ๐š ๐„๐ง๐ ๐ข๐ง๐ž๐ž๐ซ๐ข๐ง๐  ๐๐ข๐ฉ๐ž๐ฅ๐ข๐ง๐ž ๐จ๐ง ๐ƒ๐š๐ญ๐š๐›๐ซ๐ข๐œ๐ค๐ฌ (๐…๐Œ๐‚๐† ๐ƒ๐จ๐ฆ๐š๐ข๐ง)!

๐‰๐ฎ๐ฌ๐ญ ๐›๐ฎ๐ข๐ฅ๐ญ ๐š๐ง ๐„๐ง๐-๐ญ๐จ-๐„๐ง๐ ๐ƒ๐š๐ญ๐š ๐„๐ง๐ ๐ข๐ง๐ž๐ž๐ซ๐ข๐ง๐  ๐๐ข๐ฉ๐ž๐ฅ๐ข๐ง๐ž ๐จ๐ง ๐ƒ๐š๐ญ๐š๐›๐ซ๐ข๐œ๐ค๐ฌ (๐…๐Œ๐‚๐† ๐ƒ๐จ๐ฆ๐š๐ข๐ง)! When two companies merge, unifying their messy, disparat

Aug 22, 20262 min read

๐Ÿ“Œ ๐‡๐จ๐ฐ ๐ฐ๐จ๐ฎ๐ฅ๐ ๐ฒ๐จ๐ฎ ๐ก๐š๐ง๐๐ฅ๐ž ๐š ๐ฌ๐ข๐ญ๐ฎ๐š๐ญ๐ข๐จ๐ง ๐ฐ๐ก๐ž๐ซ๐ž ๐œ๐ซ๐ข๐ญ๐ข๐œ๐š๐ฅ ๐„๐“๐‹ ๐ฉ๐ข๐ฉ๐ž๐ฅ๐ข๐ง๐ž ๐ก๐š๐ฌ ๐Ÿ๐š๐ข๐ฅ๐ž๐ ๐ข๐ง ๐ฉ๐ซ๐จ๐

โœ”๏ธ According to me, these are the few steps you could take:-โœ… Access the failure by checking logs, error messages, and monitoring alerts to identify the root cause, like a schema changeโœ… Check how much data has been affected and also if the downstrea...

Aug 27, 20251 min read

Aatish's Data blog

9 posts