This course provides a comprehensive introduction to Big Data Engineering, focusing on the Apache Spark ecosystem as the core platform for distributed data processing. The curriculum covers foundational concepts, the Hadoop ecosystem, and modern data engineering practices, with extensive hands‑on experience in building scalable data pipelines using Spark, PySpark, and Spark SQL.
Target Audience: Aspiring data engineers, data analysts, software developers, and IT professionals.
Prerequisites: Basic programming (preferably Python); familiarity with SQL (SELECT, JOIN, GROUP BY); understanding of fundamental data processing concepts.
Format: 12–14 weeks (3 contact hours/week) with lectures, hands‑on labs, notebook assignments, and a capstone project.
Learning Objectives: Define Big Data, understand volume/velocity/variety; distinguish structured vs. semi‑structured vs. unstructured data; explain CAP theorem and eventual consistency; understand limitations of traditional client‑server processing.
Learning Objectives: Explain Hadoop's core components and design principles; work with HDFS commands; understand MapReduce limitations; identify broader Hadoop ecosystem tools.
Hands‑On: Basic HDFS operations (upload, download, directories); run a simple MapReduce job and analyze execution flow.
Learning Objectives: Understand the role of NoSQL in managing Big Data; differentiate key‑value, document, column‑family, and graph databases; explain schema‑less models; survey prominent NoSQL systems.
Learning Objectives: Understand Spark's architecture and advantages over MapReduce; distinguish batch vs. in‑memory processing; explain driver/executors/workers; install and configure Spark standalone.
Hands‑On: Install Apache Spark standalone; invoke the Spark shell and perform basic operations; create a SparkContext and load files.
Learning Objectives: Work with Resilient Distributed Datasets (RDDs); distinguish transformations vs. actions; understand persistence/storage levels; apply MapReduce‑style operations with Pair RDDs.
Hands‑On: Build applications using RDDs in PySpark; perform groupBy and join operations; read/write to HDFS.
Learning Objectives: Apply Spark SQL for structured data; work with DataFrames; perform aggregations, joins, window functions; query data using SQL syntax within Spark.
Hands‑On: Read CSV/JSON into DataFrames; write transformed data in Parquet; perform grouping/aggregation on e‑commerce data.
Learning Objectives: Handle complex types (arrays, maps, structs); apply relational and set operations; implement UDFs; apply performance best practices.
Hands‑On: Work with complex types in e‑commerce data; implement UDFs; optimize DataFrame operations with partitioning.
Learning Objectives: Design and implement end‑to‑end ETL pipelines; build batch processing solutions for production; integrate Spark with other Big Data tools; apply production best practices.
Hands‑On: Build a complete ETL pipeline reading from multiple sources; implement incremental load strategies; optimize with caching and partitioning.
Learning Objectives: Build streaming applications using Spark Structured Streaming; apply window aggregations; implement fault‑tolerant streaming; understand sources, output modes, and sinks.
Hands‑On: Build a streaming application; implement window aggregation; process real‑time data with Kafka integration.
Learning Objectives: Understand Lakehouse architecture; work with Delta Lake for ACID‑compliant storage; implement schema evolution and versioning; learn the Medallion architecture.
Hands‑On: Work with Delta Lake tables; implement versioning and rollback; organize data using catalogs, schemas, and volumes.
Learning Objectives: Optimize Spark job performance; monitor Spark applications; implement data quality frameworks; deploy with CI/CD and containerization.
Hands‑On: Analyze and optimize a poorly‑performing job; implement automated data quality checks; deploy a pipeline to production.
Learning Objectives: Design end‑to‑end pipelines using Spark, dbt, Airflow; integrate multiple big data technologies; apply best practices; present findings in a technical report.
Capstone Project: Design and implement a production‑ready pipeline that ingests batch + streaming data, transforms it through modular ETL, optimizes performance, and demonstrates monitoring & deployment. Deliverables: working pipeline, code repository, technical report.
By the end of this course, students will be able to:
© 2026 • Course Outline • Big Data Engineering with Spark
Welcome to Module 3 of our Big Data journey!
In Module 1, we learned that Big Data is too big for one computer. We learned about the three "V"s: Volume, Velocity, and Variety.
In Module 2, we explored Hadoop and how it stores and processes Big Data using HDFS and MapReduce. We saw how Hadoop splits files into blocks and replicates them across many computers.
Now, in Module 3, we are going to explore a very important topic: NoSQL Databases.
You might have heard of databases before. Maybe your school has a database of student records. Maybe your parents have a database of contacts on their phones. These are usually relational databases (like SQL databases). They store data in tables with rows and columns.
But Big Data is different. Big Data comes in many shapes and sizes. Some data is structured (like tables), some is semi-structured (like emails), and some is unstructured (like videos).
Relational databases are great for structured data. But they struggle with Big Data. They are slow when you have billions of records. They are not flexible enough to handle different types of data.
That is where NoSQL databases come in! NoSQL databases are designed specifically for Big Data. They are fast, flexible, and can handle all kinds of data.
In this module, we will learn:
Let's begin our adventure into the world of NoSQL!
By the end of this module, you will be able to:
Once upon a time, there was a very special library called the Magic Library. This library was different from any other library in the world.
A normal library has books arranged on shelves in a neat, organized way. Each book has a specific place. If you want a book, you go to its section and find it.
But the Magic Library was different. There were no shelves. There were no sections. There was no order at all!
Some books were on tables. Some were on the floor. Some were hanging from the ceiling. Some books were big, some were small, and some were shaped like triangles!
There were also scrolls, maps, photographs, and even jars containing tiny messages.
People were confused at first. "How can you find anything?" they asked.
The librarian, Mrs. Chioma, smiled and said, "This library is not about order. It is about flexibility. Every book has a unique address. I can find any book instantly because I know exactly where it is."
She explained: "There are four sections in this library. Each section stores things differently."
This library is exactly like NoSQL databases! They are designed to be flexible and handle all kinds of data. They store data in different ways to make it fast and easy to access.
Now, let's explore each section of the Magic Library and learn about NoSQL databases!
Definition: A database is a place where data is stored and organized so that it can be easily accessed, managed, and updated.
Why it is important: Almost every application you use uses a database. Your school uses a database to store student records. Banks use databases to store customer information. Without databases, we could not manage large amounts of information.
Simple explanation: A database is like a big box where you keep all your important information.
Real-life example: Your phone's contact list is a database. It stores names, phone numbers, and email addresses.
School example: The school office keeps a database of all students, their classes, and their grades.
Home example: Your mother has a database of recipes. Each recipe has a name, ingredients, and instructions.
Nigerian example: A bank in Nigeria has a database of all its customers and their account details.
+---------------------------------------------+
| DATABASE |
+---------------------------------------------+
| A place to store and organize data |
| Examples: Contacts, recipes, student |
| records, bank accounts |
+---------------------------------------------+
Mini summary: A database is a place to store and organize information.
Definition: A relational database (also called an SQL database) stores data in tables with rows and columns. Each table has a fixed structure.
Why it is important: Relational databases are very organized and reliable. They are used by most businesses and organizations.
Simple explanation: A relational database is like a big spreadsheet with multiple sheets (tables). Each sheet has rows and columns.
Real-life example: A school uses a relational database to store student records. There is a table for students, a table for classes, and a table for grades.
School example: The school timetable is a relational table. It has columns for time, subject, and teacher.
Home example: Your family uses a spreadsheet to track expenses. It has columns for date, item, and cost.
Nigerian example: A bank uses a relational database to store customer accounts. Each account has a unique ID, name, and balance.
+---------------------------------------------+
| RELATIONAL DATABASE |
+---------------------------------------------+
| Table: Students |
| +----------+----------+----------+ |
| | StudentID| Name | Class | |
| +----------+----------+----------+ |
| | 1001 | Ade | Primary 3| |
| | 1002 | Bola | Primary 4| |
| | 1003 | Chioma | Primary 5| |
| +----------+----------+----------+ |
+---------------------------------------------+
Mini summary: Relational databases store data in tables with rows and columns.
Definition: NoSQL stands for "Not Only SQL." It is a type of database that does not use tables with rows and columns. Instead, it stores data in different, more flexible ways.
Why it is important: NoSQL databases are designed for Big Data. They are fast, flexible, and can handle different types of data.
Simple explanation: NoSQL is a different way to store data. It does not use tables. It uses other structures like key-value pairs, documents, or graphs.
Real-life example: Facebook uses NoSQL databases to store billions of messages and photos.
School example: A school uses a NoSQL database to store student portfolios. Each portfolio has different types of work: essays, drawings, and videos.
Home example: Your family uses a NoSQL database to store a collection of recipes. Each recipe has a different structure (some have photos, some have videos, some have comments).
Nigerian example: Jumia uses NoSQL to store product information. Each product has different attributes (size, color, price, reviews).
+---------------------------------------------+
| NOSQL |
+---------------------------------------------+
| "Not Only SQL" |
| Flexible data storage |
| Designed for Big Data |
| No fixed table structure |
+---------------------------------------------+
| Types: Key-Value, Document, Column-Family, |
| Graph |
+---------------------------------------------+
Mini summary: NoSQL is a flexible type of database designed for Big Data.
Definition: NoSQL databases are better than relational databases for Big Data because they are:
Why it is important: Big Data needs databases that can grow and handle variety. NoSQL databases are perfect for this.
Simple explanation: NoSQL is like a flexible toy that can change shape. Relational is like a fixed shape that cannot change.
Real-life example: Twitter uses NoSQL to handle millions of tweets per second.
School example: A school uses NoSQL to store different types of student work (text, images, videos).
Home example: Your family uses NoSQL to store different types of memories (photos, videos, diaries).
Nigerian example: A Nigerian fintech company uses NoSQL to handle millions of transactions.
+---------------------------------------------+
| WHY NOSQL FOR BIG DATA? |
+---------------------------------------------+
| 1. Scalable (add more computers) |
| 2. Flexible (handle any data type) |
| 3. Fast (process huge data quickly) |
| 4. Distributed (runs on many computers) |
+---------------------------------------------+
Mini summary: NoSQL is perfect for Big Data because it is scalable, flexible, fast, and distributed.
Definition: There are four main types of NoSQL databases:
Why it is important: Each type is good for different use cases. You choose the type that fits your needs.
Simple explanation: There are four different ways to organize data in a NoSQL database, like four different ways to organize your toys.
Real-life example: Amazon uses key-value, Document, and graph databases for different purposes.
School example: A school might use a document database for student portfolios and a graph database for social networks.
Home example: Your family might use a key-value database for a shopping list and a graph database for the family tree.
Nigerian example: A Nigerian bank uses column-family for transaction data and graph for fraud detection.
+---------------------------------------------+
| FOUR TYPES OF NOSQL |
+---------------------------------------------+
| 1. Key-Value (like a dictionary) |
| 2. Document (like a book page) |
| 3. Column-Family (like a giant spreadsheet)|
| 4. Graph (like a spider web) |
+---------------------------------------------+
Mini summary: There are four types of NoSQL databases: key-value, document, column-family, and graph.
Definition: A key-value database stores data as a collection of key-value pairs. A key is a unique name, and a value is the data associated with that key.
Why it is important: Key-value databases are very fast. They are great for caching and storing simple data.
Simple explanation: It is like a dictionary. You look up a word (key) and you find its meaning (value).
Real-life example: A phone book is a key-value store. The key is the name, and the value is the phone number.
School example: A school uses a key-value store for student IDs. The key is the student ID, and the value is the student's name.
Home example: Your family uses a key-value store for chores. The key is the chore, and the value is the person responsible.
Nigerian example: A Nigerian e-commerce site uses a key-value store to cache product prices for quick access.
Popular key-value databases: Redis, Riak, DynamoDB.
+---------------------------------------------+
| KEY-VALUE DATABASE |
+---------------------------------------------+
| +-------------------+ |
| | Key: "Ade" | |
| | Value: "0801234" | |
| +-------------------+ |
| +-------------------+ |
| | Key: "Bola" | |
| | Value: "0805678" | |
| +-------------------+ |
| +-------------------+ |
| | Key: "Chioma" | |
| | Value: "0809012" | |
| +-------------------+ |
+---------------------------------------------+
Mini summary: Key-value databases store data as key-value pairs, like a dictionary.
Definition: A document database stores data as documents, usually in JSON format. Each document can have different fields.
Why it is important: Document databases are very flexible. You can store different types of data in the same collection.
Simple explanation: It is like a box of papers. Each paper (document) can have different information on it.
Real-life example: A library stores information about books. Each book (document) has different fields: title, author, year, pages.
School example: A school stores student portfolios. Each portfolio (document) has different fields: name, class, projects, grades.
Home example: Your family stores recipes. Each recipe (document) has different fields: name, ingredients, steps, time.
Nigerian example: Jumia stores product information. Each product (document) has different fields: name, price, description, reviews.
Popular document databases: MongoDB, CouchDB, Firebase.
+---------------------------------------------+
| DOCUMENT DATABASE |
+---------------------------------------------+
| Document 1: |
| { |
| name: "Ade", |
| class: "Primary 3", |
| subjects: ["Math", "English"] |
| } |
| |
| Document 2: |
| { |
| name: "Bola", |
| class: "Primary 4", |
| projects: ["Art", "Science"] |
| } |
+---------------------------------------------+
Mini summary: Document databases store data as flexible documents, like JSON objects.
Definition: A column-family database stores data in columns instead of rows. It is like a giant spreadsheet where each column can be different.
Why it is important: Column-family databases are great for large-scale data analysis and time-series data.
Simple explanation: It is like a spreadsheet that can grow sideways. Each row can have different columns.
Real-life example: A weather station stores temperature data. Each column is a time, and each row is a location.
School example: A school stores attendance data. Each column is a date, and each row is a student.
Home example: Your family stores monthly expenses. Each column is a month, and each row is an expense category.
Nigerian example: MTN Nigeria uses a column-family database to store call records. Each column is a time period, and each row is a user.
Popular column-family databases: Cassandra, HBase, Amazon SimpleDB.
+---------------------------------------------+
| COLUMN-FAMILY DATABASE |
+---------------------------------------------+
| Row Key: User1 |
| Column: Jan -> 200 minutes |
| Column: Feb -> 150 minutes |
| Column: Mar -> 180 minutes |
| |
| Row Key: User2 |
| Column: Jan -> 100 minutes |
| Column: Feb -> 120 minutes |
| Column: Mar -> 90 minutes |
+---------------------------------------------+
Mini summary: Column-family databases store data in columns, like a giant spreadsheet.
Definition: A graph database stores data as nodes (entities) and edges (relationships). It is like a spider web of connected data.
Why it is important: Graph databases are great for analyzing relationships between data points.
Simple explanation: It is like a family tree. Each person (node) is connected to others (edges) by relationships.
Real-life example: Facebook uses a graph database to store friendships. Each person is a node, and each friendship is an edge.
School example: A school uses a graph database to show which students are in which clubs. Students and clubs are nodes, and memberships are edges.
Home example: Your family uses a graph database to build a family tree. Each person is a node, and relationships are edges.
Nigerian example: A Nigerian bank uses a graph database to detect fraud. They connect accounts, transactions, and users.
Popular graph databases: Neo4j, Dgraph, TigerGraph.
+---------------------------------------------+
| GRAPH DATABASE |
+---------------------------------------------+
| Ade ---- friends ---- Bola |
| | | |
| friends friends |
| | | |
| Chioma ---- friends ---- David |
| |
| Nodes: Ade, Bola, Chioma, David |
| Edges: friendships |
+---------------------------------------------+
Mini summary: Graph databases store data as nodes and edges, like a spider web.
Definition: Schema-less means that the database does not require a fixed structure. You can store different types of data in the same place.
Why it is important: Schema-less databases are flexible. You can change the structure without affecting existing data.
Simple explanation: It is like a box where you can put anything, without deciding in advance what will go in.
Real-life example: A social media site stores user profiles. Some users have photos, some have videos, and some have text. All are stored in the same database.
School example: A school stores student projects. Each project is different. Some are essays, some are drawings, some are videos.
Home example: Your family stores a collection of memories. Some are photos, some are videos, some are documents.
Nigerian example: A Nigerian e-commerce site stores product information. Each product has different attributes (size, color, weight, material).
+---------------------------------------------+
| SCHEMA-LESS DATA |
+---------------------------------------------+
| Document 1: { name: "Ade", age: 10 } |
| Document 2: { name: "Bola", class: "P4" } |
| Document 3: { name: "Chioma", subjects: |
| ["Math", "Science"] } |
| |
| All documents are in the same collection! |
+---------------------------------------------+
Mini summary: Schema-less means you can store different types of data without a fixed structure.
Definition: Shared-nothing means each computer in a system has its own resources (memory, disk, CPU). They do not share anything.
Why it is important: Shared-nothing architecture allows NoSQL databases to scale easily. Each computer works independently.
Simple explanation: Each student has their own book and pencil. They do not have to share.
Real-life example: In a restaurant, each chef has their own kitchen station.
School example: Each student has their own desk and chair.
Home example: Each person in the family has their own toothbrush.
Nigerian example: Each driver has their own car. They do not share one car.
+---------------------------------------------+
| SHARED-NOTHING ARCHITECTURE |
+---------------------------------------------+
| Computer 1: own disk, own memory, own CPU |
| Computer 2: own disk, own memory, own CPU |
| Computer 3: own disk, own memory, own CPU |
| |
| They do not share anything! |
| This makes them fast and reliable. |
+---------------------------------------------+
Mini summary: Shared-nothing means each computer has its own resources.
Definition: SQL and NoSQL are two different ways to store data. They are good for different things.
Why it is important: Knowing the difference helps you choose the right database for your project.
Simple explanation: SQL is like a filing cabinet with organized folders. NoSQL is like a big box where you can put anything.
Real-life example: A bank uses SQL for customer accounts. A social media site uses NoSQL for user posts.
School example: A school uses SQL for student records. A school uses NoSQL for student portfolios.
Home example: Your family uses SQL for a budget spreadsheet. Your family uses NoSQL for a photo collection.
Nigerian example: A bank uses SQL for account balances. A fintech uses NoSQL for transaction data.
+---------------------------------------------+
| SQL vs. NOSQL |
+---------------------------------------------+
| +-------------------+-------------------+ |
| | SQL | NoSQL | |
| +-------------------+-------------------+ |
| | Tables | Various structures| |
| | Fixed schema | Schema-less | |
| | Good for relations| Good for Big Data | |
| | Vertical scaling | Horizontal scaling| |
| | ACID transactions | Eventual consistency| |
| +-------------------+-------------------+ |
+---------------------------------------------+
Mini summary: SQL is for structured data; NoSQL is for Big Data.
Definition: Cassandra is a popular NoSQL database that uses a column-family architecture. It is designed for high availability and scalability.
Why it is important: Cassandra is used by many large companies to handle huge amounts of data.
Simple explanation: Cassandra is like a giant spreadsheet that can grow sideways and handle millions of rows.
Real-life example: Netflix uses Cassandra to store user profiles and viewing history.
School example: A school uses Cassandra to store attendance records for all students across many years.
Home example: Your family uses Cassandra to store monthly bills for many years.
Nigerian example: MTN Nigeria uses Cassandra to store call records for millions of users.
+---------------------------------------------+
| CASSANDRA |
+---------------------------------------------+
| Column-family database |
| High availability |
| Scalable |
| Used by Netflix, Twitter, MTN |
+---------------------------------------------+
Mini summary: Cassandra is a popular column-family NoSQL database.
Definition: MongoDB is a popular NoSQL database that uses a document architecture. It stores data in JSON-like documents.
Why it is important: MongoDB is very flexible and easy to use. It is great for web applications.
Simple explanation: MongoDB is like a box of papers where each paper (document) can have different information.
Real-life example: eBay uses MongoDB to store product listings.
School example: A school uses MongoDB to store student portfolios with different types of work.
Home example: Your family uses MongoDB to store recipes with different formats.
Nigerian example: Jumia uses MongoDB to store product information with different attributes.
+---------------------------------------------+
| MONGODB |
+---------------------------------------------+
| Document database |
| Flexible schema |
| JSON-like documents |
| Used by eBay, Jumia, many startups |
+---------------------------------------------+
Mini summary: MongoDB is a popular document NoSQL database.
Definition: Many Nigerian companies use NoSQL databases to handle Big Data.
Why it is important: NoSQL is helping Nigerian companies grow and serve their customers better.
Simple explanation: NoSQL is used in Nigeria to handle large amounts of data quickly.
Real-life examples in Nigeria:
School example: A Nigerian school uses MongoDB to store student portfolios.
Home example: A Nigerian family uses a NoSQL database to store photos and videos.
Mini summary: NoSQL is used in Nigeria by many companies to handle Big Data.
+-------------------------------------------------+
| NOSQL DATABASES |
+-------------------------------------------------+
| |
| +----------------------------------------+ |
| | KEY-VALUE | |
| | (Key → Value) | |
| | Example: "Ade" → "0801234" | |
| +----------------------------------------+ |
| |
| +----------------------------------------+ |
| | DOCUMENT | |
| | (JSON-like documents) | |
| | Example: {name:"Ade", age:10} | |
| +----------------------------------------+ |
| |
| +----------------------------------------+ |
| | COLUMN-FAMILY | |
| | (Columns like a spreadsheet) | |
| | Example: User1: Jan 200, Feb 150 | |
| +----------------------------------------+ |
| |
| +----------------------------------------+ |
| | GRAPH | |
| | (Nodes and edges) | |
| | Example: Ade → friends → Bola | |
| +----------------------------------------+ |
+-------------------------------------------------+
+---------------------------------------------+
| KEY-VALUE DATABASE |
+---------------------------------------------+
| +-------------------+ |
| | Key: "Ade" | |
| | Value: "0801234" | |
| +-------------------+ |
| +-------------------+ |
| | Key: "Bola" | |
| | Value: "0805678" | |
| +-------------------+ |
| +-------------------+ |
| | Key: "Chioma" | |
| | Value: "0809012" | |
| +-------------------+ |
+---------------------------------------------+
+---------------------------------------------+
| DOCUMENT DATABASE |
+---------------------------------------------+
| Document 1: |
| { |
| name: "Ade", |
| class: "Primary 3", |
| subjects: ["Math", "English"] |
| } |
| |
| Document 2: |
| { |
| name: "Bola", |
| class: "Primary 4", |
| projects: ["Art", "Science"] |
| } |
+---------------------------------------------+
+---------------------------------------------+
| COLUMN-FAMILY DATABASE |
+---------------------------------------------+
| Row Key: User1 |
| Column: Jan -> 200 minutes |
| Column: Feb -> 150 minutes |
| Column: Mar -> 180 minutes |
| |
| Row Key: User2 |
| Column: Jan -> 100 minutes |
| Column: Feb -> 120 minutes |
| Column: Mar -> 90 minutes |
+---------------------------------------------+
+---------------------------------------------+
| GRAPH DATABASE |
+---------------------------------------------+
| Ade ---- friends ---- Bola |
| | | |
| friends friends |
| | | |
| Chioma ---- friends ---- David |
| |
| Nodes: Ade, Bola, Chioma, David |
| Edges: friendships |
+---------------------------------------------+
| Feature | SQL (Relational) | NoSQL |
|---|---|---|
| Structure | Tables with rows and columns | Flexible (key-value, document, etc.) |
| Schema | Fixed | Schema-less |
| Scalability | Vertical (scale-up) | Horizontal (scale-out) |
| Data type | Structured | All types |
| Best for | Transactions, relationships | Big Data, analytics |
| Examples | MySQL, PostgreSQL | MongoDB, Cassandra |
| Type | Storage | Example | Use Case |
|---|---|---|---|
| Key-Value | Key → Value | Redis, DynamoDB | Caching, session storage |
| Document | JSON-like documents | MongoDB, CouchDB | Web apps, content management |
| Column-Family | Columns | Cassandra, HBase | Time-series, analytics |
| Graph | Nodes and edges | Neo4j, Dgraph | Social networks, fraud detection |
In this module, we learned about NoSQL databases and how they are used for Big Data.
We started with a story about the Magic Library to understand the four types of NoSQL databases.
We learned that NoSQL stands for "Not Only SQL." NoSQL databases are flexible, scalable, and designed for Big Data.
We explored the four types of NoSQL databases:
We learned about schema-less data and how it allows flexibility. We also learned about shared-nothing architecture and how it makes NoSQL scalable.
We compared SQL and NoSQL and saw that each is good for different use cases.
We saw examples from Nigeria, including MTN, Jumia, and Flutterwave using NoSQL to handle Big Data.
Remember: NoSQL databases are powerful tools for Big Data. In the next module, we will learn about Apache Spark, which is a fast data processing engine that works with NoSQL databases.
Match the term on the left with the correct definition on the right.
| Term | Definition |
|---|---|
| 1. NoSQL | A. Stores data as key-value pairs |
| 2. Key-Value | B. Stores data as JSON-like documents |
| 3. Document | C. Stores data as nodes and edges |
| 4. Column-Family | D. Stores data in columns |
| 5. Graph | E. Not Only SQL |
| 6. Schema-less | F. No fixed structure |
| 7. Shared-Nothing | G. Each computer has its own resources |
Answers: 1-E, 2-A, 3-B, 4-D, 5-C, 6-F, 7-G
Activity: "Design a NoSQL Solution"
Instructions:
Activity: "My NoSQL Example"
Instructions:
Title: "NoSQL for a Nigerian Healthcare System"
Instructions:
Assignment: "Explore MongoDB"
Instructions:
Title: "Choose the Right NoSQL"
Instructions:
1-E, 2-A, 3-B, 4-D, 5-C, 6-F, 7-G
In the next module, we will learn about Apache Spark and how it is used for Big Data processing.
We will explore:
To prepare, think about these questions:
We will continue our journey into Big Data by exploring the world of Apache Spark. Get ready for an exciting adventure!
End of Module 3
Well done! You have completed the third module of Big Data Engineering with Spark.
Welcome to the world of Big Data and Distributed Computing!
Have you ever wondered how YouTube knows which videos to suggest to you? Or how your favorite online game can have millions of players at the same time without crashing? Or how a bank can check thousands of transactions every second to make sure nobody is stealing money?
All these amazing things are possible because of Big Data and Distributed Computing.
In this module, we will start from the very beginning. We will learn what "big data" really means, why regular computers sometimes cannot handle it, and how clever engineers use many computers working together to solve huge problems.
We will use simple words, fun stories, and lots of examples from everyday life — including examples from Nigeria and other places you know. By the end of this module, you will understand the basic ideas that make modern data engineering possible.
So, let's begin our adventure into the world of Big Data!
By the end of this module, you will be able to:
Once upon a time, in a busy Nigerian city, there was a huge market called Olúwo Market. Every day, thousands of people came to buy and sell goods — yams, tomatoes, shoes, phones, fabrics, and many other things.
The market leaders wanted to know exactly how many items were sold each day. They asked one person, Mr. Adebayo, to count everything.
Mr. Adebayo walked around with a small notebook and a pen. He counted one yam at a time. He counted one tomato at a time. He was very careful. But at the end of the day, he had only counted items from 10 shops. There were 500 shops in the market!
The market leaders said, "Mr. Adebayo, this is too slow! We need to know everything that was sold today — all 500 shops — before tomorrow morning."
Mr. Adebayo was worried. He could not count everything alone. He needed help.
So, he asked his friends: Chidi, Ngozi, and Fatima to help him. They each took one section of the market.
Chidi counted the food section. Ngozi counted the clothing section. Fatima counted the electronics section. And Mr. Adebayo counted the household goods section.
They all counted at the same time. When they finished, they added their numbers together. They had the total count for the whole market by evening!
This is exactly how distributed computing works. One computer (Mr. Adebayo) cannot handle big data alone. But many computers working together (Chidi, Ngozi, Fatima, and Mr. Adebayo) can finish the job quickly.
The market leaders were very happy. They said, "From now on, we will always use many people to count. And we will use many computers to handle our big data!"
And that, dear learner, is how our story introduces the idea of Big Data and Distributed Computing.
Definition: Data is any piece of information. It can be a number, a word, a picture, a sound, or anything that tells us something.
Why it is important: Everything in the world produces data. Without data, we cannot learn, decide, or improve anything.
Simple explanation: Think of data as "bits of information."
Real-life example: Your name is data. Your age is data. Your school's name is data.
School example: Your teacher writes your test score in a book. That score is data.
Home example: Your mother writes a shopping list. That list is data.
Nigerian example: The price of a bag of rice in Lagos is data. The number of people who watch Nollywood movies is data.
+-------------+
| DATA |
| |
| Any piece |
| of |
| information|
+-------------+
Mini summary: Data is just information. It is everywhere.
Definition: Big Data is data that is so large, so fast, or so complicated that normal computers cannot handle it easily.
Why it is important: We live in a world where we create huge amounts of data every second. We need special tools to manage it.
Simple explanation: Big Data is "too much information for one computer to handle."
Real-life example: Every day, people send 500 million tweets on Twitter. That is Big Data.
School example: A school has 5,000 students. If each student sends one message per day, that is 5,000 messages. That is small data. But if a school has 5 million students, that is Big Data.
Home example: Your family takes 100 photos on holiday. That is small data. But if the whole world takes 1 billion photos in one day, that is Big Data.
Nigerian example: MTN Nigeria has over 70 million subscribers. All their call records, messages, and payments make Big Data.
Small Data Big Data
---------- ---------
100 photos 1 billion photos
5,000 messages 5 million messages
1 notebook 1 billion notebooks
Mini summary: Big Data is data that is too big for one computer.
Definition: Big Data is often described using three words that start with the letter "V":
Why it is important: These three "V"s help us understand what makes data "big."
Simple explanation:
Real-life example:
School example:
Home example:
Nigerian example:
+-------------------------------------------+
| BIG DATA |
+-------------------------------------------+
| VOLUME | VELOCITY | VARIETY |
+----------+------------+-------------------+
| "How | "How fast?"| "What kind?" |
| much?" | | |
+----------+------------+-------------------+
Mini summary: Big Data has three "V"s: Volume, Velocity, and Variety.
Definition:
Why it is important: Different types of data need different tools to handle them.
Simple explanation:
Real-life example:
School example:
Home example:
Nigerian example:
+----------------+------------------+------------------+
| STRUCTURED | SEMI-STRUCTURED | UNSTRUCTURED |
+----------------+------------------+------------------+
| Clear tables | Some order | No clear order |
| Timetables | Emails | Videos |
| Spreadsheets | Report cards | Photos |
| Bank records | Letters | Voice messages |
+----------------+------------------+------------------+
Mini summary: Data can be structured, semi-structured, or unstructured.
Definition: A single computer has limits. It can only store so much data and process so much information at one time.
Why it is important: When data gets too big, we need many computers working together.
Simple explanation: One person can carry 20 kg. But if you have 1,000 kg, you need 50 people.
Real-life example: Google processes over 40,000 searches every second. One computer cannot do that.
School example: One teacher can grade 30 tests. But if the school has 1,000 tests, the school needs many teachers.
Home example: One person can wash 10 plates. But for 500 plates, you need help.
Nigerian example: One bank teller can serve 50 customers. But a bank with 5,000 customers needs many tellers and many computers.
One Computer Many Computers
------------ --------------
Can store: 1 TB Can store: 1,000 TB
Can process: 1 task/sec Can process: 1,000 tasks/sec
Works alone Works together as a team
Mini summary: One computer cannot handle Big Data. We need many computers.
Definition: Distributed computing is when many computers work together to solve a big problem.
Why it is important: Distributed computing lets us handle Big Data that a single computer cannot.
Simple explanation: It is like a team of people working together on a big project instead of one person doing it alone.
Real-life example: Google uses thousands of computers to answer your search questions.
School example: The whole class works together to clean the classroom. Each person cleans one part.
Home example: The whole family helps to cook dinner. One person chops vegetables, another cooks rice, another sets the table.
Nigerian example: A group of farmers in a village work together to harvest a big farm. Each farmer harvests one section.
+-------------------+
| Distributed |
| Computing |
+-------------------+
| Computer 1 |
| Computer 2 |
| Computer 3 |
| Computer 4 |
| Computer 5 |
+-------------------+
| All work together |
+-------------------+
Mini summary: Distributed computing means many computers working together.
Definition: Data locality means moving the work to the data, instead of moving the data to the work.
Why it is important: Moving data takes time and network bandwidth. It is faster to move the work.
Simple explanation: Instead of carrying all the food to the kitchen, you take the kitchen to the food.
Real-life example: In a factory, they bring the machines to where the materials are, not the other way around.
School example: The teacher goes to each classroom to teach, instead of all students coming to the staff room.
Home example: You take your book to the room where you want to read, instead of bringing the whole room to your book.
Nigerian example: In a market, each seller stays at their stall. The buyer goes to the seller.
Old Way New Way (Data Locality)
------- -----------------------
Move data → Work Move work → Data
(Slow) (Fast)
Mini summary: Data locality means we move the work to where the data is stored.
Definition: Shared-nothing means each computer in a system has its own memory and storage. They do not share anything.
Why it is important: If computers do not share, they do not fight over resources. They work better together.
Simple explanation: Each student has their own book and pencil. They do not have to share.
Real-life example: In a restaurant, each chef has their own kitchen station.
School example: Each student has their own desk and chair.
Home example: Each person in the family has their own toothbrush.
Nigerian example: Each driver has their own car. They do not share one car.
+-----------------------+
| Shared-Nothing |
+-----------------------+
| Computer 1: own disk, own memory |
| Computer 2: own disk, own memory |
| Computer 3: own disk, own memory |
+-----------------------+
Mini summary: Shared-nothing means each computer has its own resources.
Definition: CAP theorem says that a distributed system cannot have all three of these at the same time:
Why it is important: You have to choose which two are most important for your system.
Simple explanation: You cannot have everything. You have to choose.
Real-life example: A bank chooses consistency and partition tolerance. They can tolerate some downtime, but they must have correct records.
School example: A school must keep correct student records (Consistency) and must work even if one computer is down (Partition Tolerance). It can have a little downtime (Availability).
Home example: Your family has a shared calendar. If the internet is off, you cannot update it. That is a trade-off.
Nigerian example: A phone company must keep call records consistent and must work even if a server fails. They may have some downtime.
+-----------------------+
| CAP |
+-----------------------+
| Consistency |
| Availability |
| Partition Tolerance |
| Pick any two! |
+-----------------------+
Mini summary: CAP theorem says you can only have two out of three.
Definition:
Why it is important: Some jobs can wait; others need an answer right away.
Simple explanation: Batch is like washing all the dishes after dinner. Real-time is like washing each dish immediately.
Real-life example: A bank processes all cheques at the end of the day (batch). A credit card transaction is processed immediately (real-time).
School example: The school takes attendance at the start of the day (batch). The school bell rings to tell everyone to change classes (real-time).
Home example: You do all your homework at once (batch). You answer your phone when it rings (real-time).
Nigerian example: MTN Nigeria sends you a summary of your calls at the end of the week (batch). They show your balance immediately after you make a call (real-time).
Batch Processing Real-Time Processing
----------------- --------------------
Process later Process now
Example: End-of-day reports Example: Live updates
Slower Faster
Mini summary: Batch processing is later; real-time processing is immediate.
Definition: Scalability is the ability to grow a system to handle more data or more users.
Why it is important: A system must grow as the business grows.
Simple explanation: You can add more workers to finish a job faster.
Real-life example: Amazon adds more servers during Christmas to handle more shoppers.
School example: A school adds more classrooms when they have more students.
Home example: Your family buys a bigger refrigerator when you have more people.
Nigerian example: Jumia adds more delivery vans during sales like Black Friday.
Small System Large System
------------ ------------
1 computer 100 computers
1 server 1,000 servers
Works for 100 users Works for 1,000,000 users
Mini summary: Scalability means we can make the system bigger.
Definition:
Why it is important: Scale-out is usually cheaper and more powerful for Big Data.
Simple explanation: Scale-up is like buying a bigger truck. Scale-out is like buying more trucks.
Real-life example: Google uses scale-out. They have thousands of normal computers.
School example: Scale-up is buying bigger desks. Scale-out is buying more desks.
Home example: Scale-up is buying a bigger pot. Scale-out is buying more pots.
Nigerian example: Scale-up is one generator that is very powerful. Scale-out is many small generators.
Scale-Up Scale-Out
--------- ---------
Bigger computer More computers
More memory More machines
Faster CPU Many CPUs
More expensive per unit Cheaper per unit
Mini summary: Scale-up makes one computer bigger; scale-out adds more computers.
Definition: Eventual consistency means that, after some time, all computers in a system will have the same information.
Why it is important: In distributed systems, it is impossible to update all computers instantly. Eventual consistency allows systems to work while updates happen later.
Simple explanation: You tell your friends a secret. Not all of them hear it at the same time. But eventually, they all know the secret.
Real-life example: When you post on Facebook, not all your friends see it instantly. Eventually, they all see it.
School example: The teacher announces a change in schedule. Not all students hear it at the same time. But eventually, everyone hears it.
Home example: You tell your siblings a plan. They all learn it at different times.
Nigerian example: A market price changes. Not all traders hear about it at the same time. But eventually, everyone hears.
+----------------------------+
| Eventual Consistency |
+----------------------------+
| T1: Computer A updates |
| T2: Computer B updates |
| T3: Computer C updates |
| Eventually: All the same! |
+----------------------------+
Mini summary: Eventual consistency means updates happen over time.
Definition: Big Data skills help people understand and use large amounts of information to make better decisions.
Why it is important: The world is generating more data every day. People who can handle Big Data will have many job opportunities.
Simple explanation: If you can understand information, you can help people and businesses make good choices.
Real-life example: Companies like Google, Amazon, and Facebook hire Big Data engineers to improve their services.
School example: A school principal uses data to decide which subjects need more teachers.
Home example: Your parents use data about your test scores to know where you need help.
Nigerian example: The government uses data about farmers to decide where to build roads and provide support.
Mini summary: Big Data skills are very valuable for the future.
Definition: Nigeria has many examples of Big Data.
Why it is important: Big Data is not just in other countries. It is in Nigeria too.
Simple explanation: Many Nigerian companies use Big Data to improve their services.
Real-life examples in Nigeria:
School example: A Nigerian school can use data to track student performance and attendance.
Home example: A family uses data from their electricity meter to understand their usage.
Mini summary: Big Data is everywhere in Nigeria too.
Step 1: Big Problem
|
V
Step 2: Split into pieces
|
V
Step 3: Send to computers
|
V
Step 4: Each computer works
|
V
Step 5: Send results back
|
V
Step 6: Combine results
|
V
Step 7: Final answer
+---------------------------+
| DATA |
+---------------------------+
| |
| +----------+ +--------+|
| |Structured| |Semi- ||
| |(Tables) | |Struc- ||
| | | |tured ||
| +----------+ +--------+|
| |
| +---------------------+ |
| |Unstructured | |
| |(Videos, photos) | |
| +---------------------+ |
+---------------------------+
+-------------------+
| Master Computer |
+-------------------+
| | |
V V V
+------+ +------+ +------+
|Comp1 | |Comp2 | |Comp3 |
+------+ +------+ +------+
| | |
V V V
+-------------------+
| Combine Results |
+-------------------+
|
V
+-------------------+
| Final Answer |
+-------------------+
+---------------------------+
| CAP |
+---------------------------+
| |
| Consistency |
| / \ |
| / \ |
| / \ |
| Availability Partition |
| Tolerance |
| |
| Pick any two! |
+---------------------------+
Batch Processing Real-Time Processing
----------------- --------------------
Collect data Process data immediately
| |
V V
Wait Answer now
|
V
Process later
|
V
Get answer
| Feature | Structured Data | Unstructured Data |
|---|---|---|
| Organization | Very organized | Not organized |
| Example | Spreadsheet | Video |
| Ease of use | Easy to query | Hard to query |
| Tools | SQL, Excel | AI, machine learning |
| Feature | Scale-Up | Scale-Out |
|---|---|---|
| Meaning | Make one computer stronger | Add more computers |
| Analogy | Bigger truck | More trucks |
| Cost | Expensive | Cheaper |
| When to use | Small growth | Big growth |
| Feature | Batch | Real-Time |
|---|---|---|
| Processing time | Later | Immediate |
| Example | End-of-day reports | Bank transactions |
| Speed | Slower | Faster |
| Use case | Analysis | Live updates |
Each lesson in the "Main Lessons" section already includes a mini summary.
In this module, we learned about Big Data and distributed computing. We started with a story about Mr. Adebayo and the market to understand why we need many people to handle big jobs.
We learned that Big Data is data that is too large, too fast, or too complex for one computer. We explored the three "V"s: Volume (how much), Velocity (how fast), and Variety (what kind).
We also learned about different types of data: structured (organized in tables), semi-structured (some organization), and unstructured (no clear organization).
We discovered that one computer cannot handle Big Data. Instead, we use distributed computing, where many computers work together. We learned about data locality (moving work to data), shared-nothing architecture (each computer has its own resources), and the CAP theorem (you can only have two of Consistency, Availability, and Partition Tolerance).
We compared batch processing (processing later) and real-time processing (processing immediately). We also learned about scalability (growing the system), scale-up (making one computer stronger), and scale-out (adding more computers).
We saw examples from around the world and from Nigeria, including MTN, Jumia, and Flutterwave. We learned that Big Data skills are very valuable for the future.
Remember: Big Data is everywhere. It is changing how we live and work. By understanding Big Data, you are preparing for an exciting future!
Match the term on the left with the correct definition on the right.
| Term | Definition |
|---|---|
| 1. Big Data | A. Data that is organized in tables |
| 2. Structured Data | B. Data that is too big for one computer |
| 3. Distributed Computing | C. Many computers working together |
| 4. Batch Processing | D. Processing data later |
| 5. Real-Time Processing | E. Processing data immediately |
| 6. Scalability | F. The ability to grow the system |
| 7. Scale-Out | G. Adding more computers |
Answers: 1-B, 2-A, 3-C, 4-D, 5-E, 6-F, 7-G
Activity: "The Big Data Team Challenge"
Instructions:
Activity: "My Data Collection"
Instructions:
Title: "Design a Big Data System for a Nigerian Market"
Instructions:
Assignment: "Data Collection and Organization"
Instructions:
Title: "The CAP Theorem Challenge"
Instructions:
1-B, 2-A, 3-C, 4-D, 5-E, 6-F, 7-G
In the next module, we will learn about the Hadoop Ecosystem and Distributed File Storage. We will explore how data is stored across many computers using the Hadoop Distributed File System (HDFS). We will also learn about the MapReduce programming model and how it processes large datasets.
To prepare, think about these questions:
We will continue our journey into Big Data by looking at how storage and processing work in a distributed environment. Get ready for an exciting adventure!
End of Module 1
Well done! You have completed the first module of Big Data Engineering with Spark.
Welcome to Module 2 of our Big Data journey!
In Module 1, we learned that Big Data is too big for one computer. We discovered that we need many computers working together (distributed computing) to handle Big Data. We also learned about the three "V"s: Volume, Velocity, and Variety.
Now, in Module 2, we are going to look at a real system that was created to handle Big Data. That system is called Hadoop.
Hadoop is like a giant team of computers that work together to store and process huge amounts of information. Think of it as a super-powered version of the market counting team we met in Module 1!
In this module, we will learn:
Let's begin our adventure into the Hadoop ecosystem!
By the end of this module, you will be able to:
Imagine a very, very big city called Data City. In this city, there is a huge library called the Big Data Library. This library has billions of books, and every day, thousands of new books arrive.
There is only one librarian, Mr. Olu. He is very smart and hardworking. But there is just too much for one person.
Mr. Olu cannot find a book quickly because there are too many. He cannot put all the new books on the shelves because there is no space. He is overwhelmed.
The city leaders decided to solve this problem. They built a new library called the Hadoop Library.
In the Hadoop Library, there are hundreds of librarians. Each librarian is responsible for a small section of the library. When a book arrives, the head librarian sends it to one of the librarians. That librarian takes care of the book and puts it on a shelf.
If a book is very popular, the head librarian might make copies of the book and send the copies to other librarians. That way, if one librarian is busy, another librarian can help a reader find the book quickly.
When a reader comes to the library to ask a question, the head librarian gives the question to all the librarians. Each librarian looks at their own section and gives part of the answer. The head librarian collects all the answers and gives the complete answer to the reader.
This is exactly how Hadoop works:
This story helps us understand the basics of Hadoop. Now, let's explore it in more detail!
Definition: Hadoop is an open-source software framework that allows us to store and process big data using many computers working together.
Why it is important: Before Hadoop, it was very hard and expensive to handle Big Data. Hadoop made it cheaper and easier for companies to work with huge amounts of information.
Simple explanation: Hadoop is a team of computers that work together to store and process very large files.
Real-life example: Facebook uses Hadoop to store and process billions of photos and messages.
School example: A school has many teachers. Each teacher takes care of a subject. Together, they educate all the students. Hadoop is like that.
Home example: Your family has a storage room. Everyone in the family puts things in the storage room. Hadoop is like a giant storage room for data.
Nigerian example: A big bank in Nigeria has many branches. Each branch handles customers in its area. Hadoop is like a system that connects all the branches to work together.
+-----------------------------------------------+
| HADOOP |
+-----------------------------------------------+
| A team of computers working together |
| to store and process Big Data |
+-----------------------------------------------+
| Open source – free to use |
| Reliable – can handle failures |
| Scalable – can grow by adding more computers |
+-----------------------------------------------+
Mini summary: Hadoop is a free, reliable, scalable system for storing and processing Big Data.
Definition: Hadoop has two main parts:
Why it is important: These two parts work together to make Hadoop powerful. HDFS stores data, and MapReduce processes it.
Simple explanation: HDFS is like a giant warehouse. MapReduce is like the workers who sort and organize items in the warehouse.
Real-life example: In a big factory, there is a warehouse (HDFS) and workers (MapReduce) who move things around.
School example: A school has a library (HDFS) and students who use the library to do research (MapReduce).
Home example: Your home has a kitchen (HDFS) and a chef (MapReduce) who prepares food.
Nigerian example: A market has stalls (HDFS) and sellers who sell things (MapReduce).
+-----------------------+
| HADOOP |
+-----------------------+
| |
| +-------------+ |
| | HDFS | |
| | (Storage) | |
| +-------------+ |
| |
| +-------------+ |
| | MapReduce | |
| | (Processing)| |
| +-------------+ |
| |
+-----------------------+
Mini summary: Hadoop = HDFS (storage) + MapReduce (processing).
Definition: HDFS (Hadoop Distributed File System) is a system that stores very large files across many computers.
Why it is important: HDFS is designed to store huge amounts of data reliably, even if some computers fail.
Simple explanation: HDFS splits a big file into smaller pieces and stores them on different computers.
Real-life example: Google uses a system like HDFS to store all the websites it has indexed.
School example: A school splits a big book into chapters and gives each chapter to a different student to read.
Home example: Your family splits a big cleaning task into different rooms. Each person cleans one room.
Nigerian example: A farm divides its land into sections. Each farmer takes care of one section.
+---------------------------------------------+
| HDFS |
+---------------------------------------------+
| Big file |
| | |
| V |
| Split into blocks |
| | |
| V |
| Block 1 -> Computer 1 |
| Block 2 -> Computer 2 |
| Block 3 -> Computer 3 |
| Block 4 -> Computer 4 |
| Block 5 -> Computer 5 |
+---------------------------------------------+
Mini summary: HDFS splits big files into blocks and stores them on different computers.
Definition:
Why it is important: Blocks make it easy to split files. Replication ensures that data is safe even if a computer fails.
Simple explanation:
Real-life example: Netflix uses block storage and replication to stream movies to millions of people.
School example: A teacher splits a test paper into pages. Each page is a block. The teacher makes copies of each page and gives them to different students to mark.
Home example: Your family splits a big puzzle into pieces. Each piece is a block. You make copies of each piece in case you lose one.
Nigerian example: A bank stores customer records across different branch offices. Each branch has a copy (replication) of the records.
+----------------------------------------------+
| HDFS Blocks and Replication |
+----------------------------------------------+
| File: movie.mp4 (1000 MB) |
| | |
| V |
| Block 1 (128 MB) -> Replica 1, 2, 3 |
| Block 2 (128 MB) -> Replica 1, 2, 3 |
| Block 3 (128 MB) -> Replica 1, 2, 3 |
| Block 4 (128 MB) -> Replica 1, 2, 3 |
| Block 5 (128 MB) -> Replica 1, 2, 3 |
+----------------------------------------------+
Mini summary: HDFS splits files into blocks and makes copies (replicas) of each block.
Definition:
Why it is important: The NameNode tells you where your data is. DataNodes store the data.
Simple explanation: NameNode is like the library's catalog. DataNodes are like the bookshelves.
Real-life example: In a large company, the receptionist (NameNode) knows where every employee sits (DataNodes).
School example: The principal (NameNode) knows which class each teacher (DataNode) is in.
Home example: Your mother (NameNode) knows where everything in the house is stored (DataNodes).
Nigerian example: The market master (NameNode) knows which stall (DataNode) sells which goods.
+---------------------------------------------+
| HDFS Architecture |
+---------------------------------------------+
| +---------+ |
| |NameNode | (Master) |
| |(Catalog)| |
| +---------+ |
| | |
| | (Tells where data is) |
| V |
| +---------+ +---------+ +---------+ |
| |DataNode | |DataNode | |DataNode | |
| |(Worker) | |(Worker) | |(Worker) | |
| +---------+ +---------+ +---------+ |
| (Slaves) |
+---------------------------------------------+
Mini summary: NameNode knows where data is; DataNodes store the data.
Definition: MapReduce is a programming model for processing large datasets in a distributed manner.
Why it is important: MapReduce allows us to process huge files quickly by splitting the work across many computers.
Simple explanation: MapReduce has two steps:
Real-life example: A census: People fill out forms (Map), and then the government combines all the forms (Reduce).
School example: Each student in a class solves a math problem (Map). The teacher collects all the answers (Reduce).
Home example: Each family member cleans their own room (Map). At the end, the whole house is clean (Reduce).
Nigerian example: In a market, each trader counts their own sales (Map). The market master adds all sales together (Reduce).
+---------------------------------------------+
| MapReduce Process |
+---------------------------------------------+
| Big Problem |
| | |
| V |
| Split into pieces |
| | |
| V |
| MAP: Each computer processes its piece |
| | |
| V |
| Combine results (SHUFFLE) |
| | |
| V |
| REDUCE: Combine all results |
| | |
| V |
| Final answer |
+---------------------------------------------+
Mini summary: MapReduce = Map (split work) + Reduce (combine results).
Definition: A classic MapReduce example is counting the number of times each word appears in a large book.
Why it is important: This example shows how MapReduce works in a simple way.
Simple explanation:
Real-life example: A newspaper wants to know which words are used most often.
School example: Students count how many times each letter appears in a paragraph.
Home example: Your family counts how many times each vegetable appears in your shopping list.
Nigerian example: A radio station counts how many times each song is requested.
+---------------------------------------------+
| Word Count with MapReduce |
+---------------------------------------------+
| Input: "To be or not to be" |
| |
| MAP: |
| Computer 1: to:1, be:1, or:1 |
| Computer 2: not:1, to:1, be:1 |
| |
| REDUCE: |
| to: 2, be: 2, or: 1, not: 1 |
+---------------------------------------------+
Mini summary: Word count is a classic MapReduce example.
Definition: Between the Map and Reduce steps, there is a stage called shuffle and sort. This is where the system organizes the output from the Map tasks and sends it to the Reduce tasks.
Why it is important: Shuffle and sort ensure that all the data for each key goes to the same Reduce task.
Simple explanation: After each computer finishes its work, it sends its results to a central place. The results are organized and sent to the right computer for the Reduce step.
Real-life example: In a school, students hand in their tests to the teacher. The teacher sorts the tests by class and then gives them to the right class teacher.
School example: The teacher collects homework from all students, sorts them by subject, and sends them to the subject teachers.
Home example: Your parents collect all the family members' ideas for dinner. They organize them by type of food before deciding.
Nigerian example: In a market, sellers give their sales reports to the market master. The market master organizes the reports by product type.
+---------------------------------------------+
| Shuffle and Sort |
+---------------------------------------------+
| MAP 1: (apple, 3), (banana, 2) |
| MAP 2: (banana, 1), (apple, 1) |
| MAP 3: (banana, 4), (orange, 1) |
| |
| SHUFFLE (organize by key): |
| apple: [(3), (1)] |
| banana: [(2), (1), (4)] |
| orange: [(1)] |
| |
| REDUCE: |
| apple: 4 |
| banana: 7 |
| orange: 1 |
+---------------------------------------------+
Mini summary: Shuffle and sort organize data between Map and Reduce.
Definition: MapReduce is powerful, but it has some problems.
Why it is important: Understanding the limitations helps us see why Spark was created.
Simple explanation: MapReduce has these problems:
Real-life example: Imagine a bus that stops at every stop. It is slow. MapReduce is like that bus.
School example: A teacher who only checks homework once a week. That is batch processing.
Home example: Waiting until the end of the week to wash all the dishes.
Nigerian example: A business that only counts its sales at the end of the month.
+---------------------------------------------+
| MapReduce Limitations |
+---------------------------------------------+
| 1. Slow (writes to disk) |
| 2. Batch only (no real-time) |
| 3. Hard to write code |
| 4. Not good for machine learning |
| 5. No in-memory processing |
+---------------------------------------------+
Mini summary: MapReduce is slow, batch-only, and hard to use.
Definition:
Why it is important: Different jobs need different processing.
Simple explanation:
Real-life example:
School example:
Home example:
Nigerian example:
Batch Processing Real-Time Processing
----------------- --------------------
Process later Process immediately
Example: Reports Example: Transactions
Slower Faster
Used for analysis Used for live systems
Mini summary: Batch is later; real-time is immediate.
Definition: The Hadoop ecosystem is a collection of tools that work with Hadoop to perform different tasks.
Why it is important: Hadoop alone is not enough. The ecosystem provides tools for storage, processing, analytics, and more.
Simple explanation: Hadoop is like a smartphone. The ecosystem is like all the apps you can install on it.
Real-life example: Google has many tools like Gmail, Google Drive, and Google Maps. Hadoop has similar tools like Hive, Pig, and HBase.
School example: A school has a library (Hadoop). The ecosystem is like the textbooks, workbooks, and other resources in the library.
Home example: Your home has a kitchen (Hadoop). The ecosystem is like the pots, pans, knives, and other utensils.
Nigerian example: A bank has a main office (Hadoop). The ecosystem is like the branches, ATMs, and mobile apps that support it.
+---------------------------------------------+
| Hadoop Ecosystem |
+---------------------------------------------+
| +---------+ +---------+ +---------+ |
| | HDFS | |MapReduce| | YARN | |
| +---------+ +---------+ +---------+ |
| |
| +---------+ +---------+ +---------+ |
| | Hive | | Pig | | HBase | |
| +---------+ +---------+ +---------+ |
| |
| +---------+ +---------+ +---------+ |
| | Spark | | Flink | | Kafka | |
| +---------+ +---------+ +---------+ |
+---------------------------------------------+
Mini summary: The Hadoop ecosystem provides many tools to work with Big Data.
Definition: Hive is a tool that lets you use SQL-like queries on data stored in HDFS.
Why it is important: Many people know SQL. Hive allows them to work with Big Data using familiar queries.
Simple explanation: Hive is like a translator. It turns SQL queries into MapReduce jobs.
Real-life example: A business analyst can write SQL queries to analyze data without knowing MapReduce.
School example: Students write in English, and a translator converts it to Yoruba.
Home example: You ask your mom in English, and she translates to your native language.
Nigerian example: A bank employee writes a SQL query to get customer information without writing MapReduce code.
+---------------------------------------------+
| Hive |
+---------------------------------------------+
| SQL Query |
| | |
| V |
| Hive translates to MapReduce |
| | |
| V |
| MapReduce processes data |
| | |
| V |
| Results are returned |
+---------------------------------------------+
Mini summary: Hive lets you use SQL with Hadoop.
Definition: Pig is a tool that uses a language called Pig Latin to process data in Hadoop.
Why it is important: Pig makes it easier to write data processing scripts compared to MapReduce.
Simple explanation: Pig is like a simpler way to tell Hadoop what to do.
Real-life example: Instead of writing a long essay (MapReduce), you write a short poem (Pig Latin).
School example: Instead of writing a long report, you write bullet points.
Home example: Instead of giving long instructions, you give a simple recipe.
Nigerian example: A data analyst uses Pig to clean and transform data before analysis.
+---------------------------------------------+
| Pig |
+---------------------------------------------+
| Pig Latin Script |
| | |
| V |
| Pig translates to MapReduce |
| | |
| V |
| MapReduce processes data |
| | |
| V |
| Results are returned |
+---------------------------------------------+
Mini summary: Pig is a scripting language for Hadoop.
Definition: HBase is a NoSQL database that runs on top of HDFS. It allows real-time read and write access to large datasets.
Why it is important: HDFS is great for storing files, but it is slow for reading and writing small pieces of data. HBase solves this problem.
Simple explanation: HBase is like a giant, fast filing cabinet that sits on top of the HDFS warehouse.
Real-life example: Facebook uses HBase to store messages and user profiles.
School example: A school uses a database to find a student's record quickly.
Home example: Your family uses a phonebook to find a phone number quickly.
Nigerian example: A bank uses HBase to find customer details quickly during a transaction.
+---------------------------------------------+
| HBase |
+---------------------------------------------+
| +---------+ +---------+ +---------+ |
| | HBase | | HBase | | HBase | |
| | Table 1 | | Table 2 | | Table 3 | |
| +---------+ +---------+ +---------+ |
| | | | |
| V V V |
| +-------------------------------------+ |
| | HDFS | |
| +-------------------------------------+ |
+---------------------------------------------+
Mini summary: HBase is a NoSQL database on top of HDFS.
Definition: Hadoop is used in many industries, including telecommunications, banking, and agriculture in Nigeria.
Why it is important: Hadoop helps Nigerian companies handle large amounts of data and improve their services.
Simple explanation: Hadoop is a tool that helps Nigerian businesses grow and serve their customers better.
Real-life examples in Nigeria:
School example: A Nigerian school uses Hadoop to track student performance across different subjects.
Home example: A family uses Hadoop to track their electricity usage and save money.
Mini summary: Hadoop is used in Nigeria to improve services and grow businesses.
Step 1: Big file
|
V
Step 2: NameNode decides to store it
|
V
Step 3: Split into blocks
|
V
Step 4: Send blocks to DataNodes
|
V
Step 5: Replicate blocks
|
V
Step 6: NameNode records location
|
V
Step 7: NameNode tells you where blocks are
|
V
Step 8: You read the file
+-------------------------------------------------+
| HADOOP |
+-------------------------------------------------+
| |
| +-----------------------------------------+ |
| | NameNode (Master) | |
| | (Keeps track of where data is) | |
| +-----------------------------------------+ |
| | | | |
| V V V |
| +---------+ +---------+ +---------+ |
| |DataNode | |DataNode | |DataNode | |
| |(Worker) | |(Worker) | |(Worker) | |
| +---------+ +---------+ +---------+ |
| |
| +-----------------------------------------+ |
| | MapReduce | |
| | (Processes data in parallel) | |
| +-----------------------------------------+ |
+-------------------------------------------------+
Input Data
|
V
+-------------+
| MAP |
| (Process |
| pieces) |
+-------------+
|
V
+-------------+
| SHUFFLE |
| AND SORT |
+-------------+
|
V
+-------------+
| REDUCE |
| (Combine) |
+-------------+
|
V
Output Data
+---------------------------------------------+
| HDFS |
+---------------------------------------------+
| File: "bigdata.txt" (300 MB) |
| | |
| V |
| Block 1 (128 MB) |
| | |
| V |
| DataNode 1, DataNode 2, DataNode 3 |
| |
| Block 2 (128 MB) |
| | |
| V |
| DataNode 2, DataNode 3, DataNode 4 |
| |
| Block 3 (44 MB) |
| | |
| V |
| DataNode 3, DataNode 4, DataNode 5 |
+---------------------------------------------+
| Feature | HDFS | Regular File System |
|---|---|---|
| Storage | Distributed across many computers | Stored on one computer |
| Fault tolerance | High (replication) | Low (no replication) |
| File size | Large (GB, TB, PB) | Small (KB, MB) |
| Speed | Fast for large files | Fast for small files |
| Access | Write once, read many | Read and write multiple times |
| Feature | MapReduce | Traditional Processing |
|---|---|---|
| Processing | Distributed | Single computer |
| Speed | Slow (writes to disk) | Fast (in-memory) |
| Data size | Very large | Small to medium |
| Programming | Map and Reduce functions | Procedural or OOP |
| Fault tolerance | Yes | No |
| Feature | Batch Processing | Real-Time Processing |
|---|---|---|
| Processing time | Later | Immediate |
| Example | Monthly reports | Bank transactions |
| Latency | High | Low |
| Data size | Large | Small |
| Tool | Hadoop MapReduce | Spark, Flink |
In this module, we learned about Hadoop and the Hadoop ecosystem.
We started with a story about the Big Data Library and how many librarians (computers) work together to store and process data.
We learned that Hadoop is a framework that allows us to store and process Big Data using many computers. Hadoop has two main parts: HDFS (storage) and MapReduce (processing).
We explored how HDFS stores data by splitting files into blocks and making replicas of each block. We learned about the NameNode (the catalog) and DataNodes (the workers).
We also learned about MapReduce, which has two steps: Map (processing pieces) and Reduce (combining results). The shuffle and sort step organizes data between Map and Reduce.
We discussed the limitations of MapReduce: it is slow, batch-only, and hard to use. This will lead us to Spark in later modules.
We explored the Hadoop ecosystem, including tools like Hive (SQL), Pig (scripting), and HBase (NoSQL).
We also saw examples from Nigeria, including MTN, Jumia, and Flutterwave using Hadoop to improve their services.
Remember: Hadoop is a powerful tool for Big Data, but it has limitations. In the next module, we will learn about a faster, more modern tool: Apache Spark.
Match the term on the left with the correct definition on the right.
| Term | Definition |
|---|---|
| 1. Hadoop | A. The storage part of Hadoop |
| 2. HDFS | B. The processing part of Hadoop |
| 3. MapReduce | C. A framework for Big Data |
| 4. NameNode | D. A small piece of a file |
| 5. DataNode | E. Making copies of data |
| 6. Block | F. The master computer |
| 7. Replication | G. A worker computer |
Answers: 1-C, 2-A, 3-B, 4-F, 5-G, 6-D, 7-E
Activity: "Design a Hadoop Cluster"
Instructions:
Activity: "My MapReduce Example"
Instructions:
Title: "Hadoop for a Nigerian E-Commerce Company"
Instructions:
Assignment: "HDFS Operations"
Instructions:
Title: "Improving MapReduce"
Instructions:
1-C, 2-A, 3-B, 4-F, 5-G, 6-D, 7-E
In the next module, we will learn about NoSQL Databases for Big Data. We will explore different types of NoSQL databases, including key-value stores, document databases, and column-family stores.
To prepare, think about these questions:
We will continue our journey into Big Data by exploring how data is stored in a more flexible way. Get ready for an exciting adventure!
End of Module 2
Well done! You have completed the second module of Big Data Engineering with Spark.
Welcome to Module 4 of our Big Data journey!
In Module 1, we learned that Big Data is too big for one computer and needs many computers working together.
In Module 2, we explored Hadoop and its two main parts: HDFS (storage) and MapReduce (processing). We learned that MapReduce is powerful but slow because it writes data to disk between steps.
In Module 3, we explored NoSQL databases and how they store Big Data in flexible ways.
Now, in Module 4, we are going to meet the hero of Big Data processing: Apache Spark!
Spark is a game-changer. It is much faster than MapReduce because it processes data in memory (RAM) instead of writing to disk all the time. It is like the difference between a cheetah and a turtle!
In this module, we will learn:
Let's begin our adventure into the world of Apache Spark!
By the end of this module, you will be able to:
Once upon a time, there was a delivery company called MapReduce Delivery. They delivered packages all over the city.
This is how they worked: Every day, a truck would go to the warehouse, pick up packages, and deliver them. But the truck would only go once a day. And every time the truck came back, the drivers had to write everything down on paper (write to disk) before they could sort the packages for the next trip.
This system was okay, but it was slow. Customers complained that their packages took too long.
Then, a new company arrived: Spark Express.
Spark Express worked differently. They had many small, fast vans instead of one big truck. These vans could go out anytime, not just once a day. And instead of writing everything down on paper, they kept the delivery information in their memory (like a smart tablet).
When a package arrived, the van driver could immediately see where it needed to go. They could deliver packages 10 times faster than MapReduce Delivery!
Customers were very happy. They started using Spark Express for all their deliveries.
This story is exactly how Apache Spark works compared to Hadoop MapReduce:
Now, let's explore Spark in more detail!
Definition: Apache Spark is a fast, distributed data processing engine designed for Big Data. It can process data in memory (RAM), which makes it much faster than MapReduce.
Why it is important: Spark is one of the most popular Big Data tools in the world. It is used by thousands of companies to process huge amounts of data quickly.
Simple explanation: Spark is a super-fast system that helps computers work together to process Big Data.
Real-life example: Uber uses Spark to analyze ride data and optimize routes in real-time.
School example: A school uses Spark to process all students' grades and attendance records quickly.
Home example: Your family uses Spark to analyze all the photos and videos you have taken over the years.
Nigerian example: MTN Nigeria uses Spark to analyze call records and improve network coverage.
+---------------------------------------------+
| APACHE SPARK |
+---------------------------------------------+
| Fast, distributed data processing engine |
| Processes data in memory (RAM) |
| 10–100x faster than MapReduce |
| Used by thousands of companies |
+---------------------------------------------+
Mini summary: Apache Spark is a fast, in-memory data processing engine for Big Data.
Definition: Spark was created to solve the problems of MapReduce. MapReduce was too slow and hard to use for many Big Data tasks.
Why it is important: Spark made Big Data processing faster, easier, and more powerful.
Simple explanation: People needed a faster way to process Big Data, so they created Spark.
Real-life example: Google and Facebook needed to process huge amounts of data quickly. MapReduce was too slow, so they started using Spark.
School example: A school used a slow computer to calculate grades. They bought a faster computer to do it quickly.
Home example: Your family used a slow internet connection. You upgraded to a faster one.
Nigerian example: A Nigerian fintech company used MapReduce for fraud detection. It was too slow. They switched to Spark for faster detection.
+---------------------------------------------+
| WHY WAS SPARK CREATED? |
+---------------------------------------------+
| 1. MapReduce was too slow |
| 2. MapReduce was hard to use |
| 3. MapReduce was only for batch processing |
| 4. Spark is faster, easier, and more |
| powerful |
+---------------------------------------------+
Mini summary: Spark was created to overcome the limitations of MapReduce.
Definition: Spark has a distributed architecture with three main roles:
Why it is important: Understanding the architecture helps us understand how Spark works.
Simple explanation: Spark is like a team. The driver is the team leader. The executors are the team members. The cluster manager is the person who assigns tasks.
Real-life example: A construction project: the manager (driver) plans the project. The workers (executors) do the work. The project manager (cluster manager) assigns workers.
School example: A class project: the teacher (driver) gives instructions. Students (executors) do the work. The class monitor (cluster manager) helps organize.
Home example: A family dinner: the parent (driver) plans the meal. Family members (executors) cook different parts.
Nigerian example: A market: the market master (driver) coordinates. The traders (executors) sell goods.
+---------------------------------------------+
| SPARK ARCHITECTURE |
+---------------------------------------------+
| +---------------------------------------+ |
| | DRIVER | |
| | (Boss – plans the whole job) | |
| +---------------------------------------+ |
| | | | |
| V V V |
| +---------+ +---------+ +---------+ |
| |Executor | |Executor | |Executor | |
| |(Worker) | |(Worker) | |(Worker) | |
| +---------+ +---------+ +---------+ |
| |
| +---------------------------------------+ |
| | CLUSTER MANAGER | |
| | (Assigns tasks to workers) | |
| +---------------------------------------+ |
+---------------------------------------------+
Mini summary: Spark has a driver (boss), executors (workers), and a cluster manager (assigner).
Definition: The driver is the main program that runs the Spark application. It creates the SparkContext, which is the connection to Spark.
Why it is important: The driver is the brain of the Spark application. Without it, nothing works.
Simple explanation: The driver is like the captain of a ship. It tells everyone what to do.
Real-life example: A movie director (driver) tells actors (executors) what to do.
School example: The school principal (driver) gives instructions to teachers (executors).
Home example: Your mother (driver) tells the family what to do for the day.
Nigerian example: The Oba (king) (driver) gives instructions to the village heads (executors).
+---------------------------------------------+
| THE DRIVER |
+---------------------------------------------+
| The main program |
| Creates the SparkContext |
| Plans the entire job |
| Tells executors what to do |
| Collects results |
+---------------------------------------------+
Mini summary: The driver is the boss that plans and manages the Spark job.
Definition: Executors are the worker processes that run tasks on the data. They do the actual work of processing data.
Why it is important: Without executors, the driver's instructions would not be carried out.
Simple explanation: Executors are the workers who do the actual job.
Real-life example: Construction workers (executors) build the building.
School example: Students (executors) do the homework assigned by the teacher.
Home example: Family members (executors) do the chores assigned by the parent.
Nigerian example: Farmers (executors) work on the farms assigned by the chief.
+---------------------------------------------+
| EXECUTORS |
+---------------------------------------------+
| Worker processes |
| Run tasks on data |
| Store data in memory |
| Send results back to driver |
| Can be many on different computers |
+---------------------------------------------+
Mini summary: Executors are the workers that process data in Spark.
Definition: Spark has several components that work together:
Why it is important: Different components are used for different tasks.
Simple explanation: Spark is like a Swiss Army knife. It has different tools for different jobs.
Real-life example: A kitchen has different tools for different tasks (knife for cutting, pan for frying).
School example: A school has different departments (math, science, art).
Home example: Your home has different rooms for different activities.
Nigerian example: A market has different sections for different goods.
+---------------------------------------------+
| SPARK COMPONENTS |
+---------------------------------------------+
| +---------------------------------------+ |
| | Spark Core (Foundation) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | Spark SQL (Structured data) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | Spark Streaming (Real-time) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | MLlib (Machine Learning) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | GraphX (Graph processing) | |
| +---------------------------------------+ |
+---------------------------------------------+
Mini summary: Spark has different components for different tasks.
Definition: Spark and Hadoop MapReduce are both Big Data processing tools, but they are different.
Why it is important: Understanding the differences helps us choose the right tool.
Simple explanation: Spark is like a cheetah; MapReduce is like a turtle. Spark is much faster.
Real-life example: A race: Spark finishes in 1 minute; MapReduce finishes in 10 minutes.
School example: Spark is like a fast car; MapReduce is like a bicycle.
Home example: Spark is like a microwave; MapReduce is like an oven.
Nigerian example: Spark is like a fast internet connection; MapReduce is like a slow one.
| Feature | Spark | MapReduce |
|---|---|---|
| Speed | Fast (in-memory) | Slow (disk-based) |
| Processing | In-memory | Disk-based |
| Use case | Batch, streaming, ML | Batch only |
| Ease of use | Easier (high-level APIs) | Harder (low-level) |
| Iterative jobs | Excellent | Poor |
Mini summary: Spark is faster and more versatile than MapReduce.
Definition: Spark supports four programming languages:
Why it is important: You can use the language you are most comfortable with.
Simple explanation: Spark is like a restaurant that serves four different cuisines. You can choose your favorite.
Real-life example: A school teaches in four different languages.
School example: Students can choose between English, French, Yoruba, or Igbo.
Home example: Your family has four different TV channels.
Nigerian example: A market has four different sections for different goods.
+---------------------------------------------+
| LANGUAGES SUPPORTED BY SPARK |
+---------------------------------------------+
| +---------+ +---------+ +---------+ |
| | Python | | Scala | | Java | |
| | (PySpark)| | (Native)| | | |
| +---------+ +---------+ +---------+ |
| +---------+ |
| | R | |
| | | |
| +---------+ |
+---------------------------------------------+
Mini summary: Spark supports Python, Scala, Java, and R.
Definition: Spark can be installed and run in different modes:
Why it is important: You can start with local mode to learn and then move to a cluster.
Simple explanation: Installing Spark is like downloading a game. You can play it on your computer or on a bigger server.
Real-life example: You can watch a movie on your phone (local) or on a big screen (cluster).
School example: You can study in your room (local) or in the library (cluster).
Home example: You can play a game on your tablet (local) or on the big TV (cluster).
Nigerian example: You can sell goods from a small stall (local) or from a big market (cluster).
+---------------------------------------------+
| SPARK INSTALLATION MODES |
+---------------------------------------------+
| Local: On your computer (learning) |
| Standalone: On a cluster |
| YARN: On Hadoop YARN |
| Mesos: On Apache Mesos |
+---------------------------------------------+
| Start with local mode to learn! |
+---------------------------------------------+
Mini summary: Spark can run locally or on a cluster.
Definition: The Spark shell is an interactive environment where you can run Spark commands and see results immediately.
Why it is important: The shell is the best way to learn Spark. You can experiment and see what happens.
Simple explanation: The Spark shell is like a playground where you can try things out.
Real-life example: A scientist uses a lab to experiment.
School example: A student uses a notebook to practice.
Home example: You use a sketchpad to draw.
Nigerian example: A chef uses a test kitchen to try new recipes.
+---------------------------------------------+
| SPARK SHELL |
+---------------------------------------------+
| Interactive environment |
| Run Spark commands |
| See results immediately |
| Best way to learn Spark |
| Command: spark-shell or pyspark |
+---------------------------------------------+
Mini summary: The Spark shell is the best way to start learning Spark.
Definition: Functional programming is a way of writing code where you focus on functions (like mathematical functions) instead of changing data directly.
Why it is important: Spark is built on functional programming concepts. Understanding this makes Spark easier to use.
Simple explanation: Instead of telling the computer "change this," you tell it "create a new version of this."
Real-life example: Instead of editing a photo directly, you make a copy and edit the copy.
School example: Instead of writing on the original paper, you make a copy and write on the copy.
Home example: Instead of changing the original recipe, you write a new version.
Nigerian example: Instead of changing the original plan, you make a new plan.
+---------------------------------------------+
| FUNCTIONAL PROGRAMMING |
+---------------------------------------------+
| Focus on functions |
| Do not change data directly |
| Create new data instead |
| Used by Spark |
| Examples: map, filter, reduce |
+---------------------------------------------+
Mini summary: Functional programming is about using functions to create new data, not changing existing data.
Definition: Spark uses functional programming because it is easier to distribute work across many computers.
Why it is important: Functional programming makes it safer and easier to run code on many computers.
Simple explanation: Since you don't change data directly, there is no confusion when many computers work on the same data.
Real-life example: Ten people can work on the same project if they all have their own copies and don't change the original.
School example: Ten students can each draw their own picture without changing the original.
Home example: Family members can each have their own copy of a recipe.
Nigerian example: Ten traders can each have their own stock without changing the central supply.
+---------------------------------------------+
| WHY SPARK USES FUNCTIONAL PROGRAMMING |
+---------------------------------------------+
| 1. Easier to distribute work |
| 2. Safer (no unexpected changes) |
| 3. Faster (parallel processing) |
| 4. Cleaner code |
| 5. Works well with Big Data |
+---------------------------------------------+
Mini summary: Functional programming makes Spark fast and safe for Big Data.
Definition: Map is a functional operation that applies a function to every element in a collection and returns a new collection.
Why it is important: Map is one of the most common operations in Spark.
Simple explanation: Map is like a machine that takes one thing and turns it into another thing.
Real-life example: A machine that takes apples and turns them into apple juice.
School example: A teacher takes each student's test score and adds 10 points.
Home example: You take each vegetable and chop it.
Nigerian example: A market seller takes each item and adds a price tag.
+---------------------------------------------+
| MAP |
+---------------------------------------------+
| Input: [1, 2, 3, 4] |
| Function: x -> x * 2 |
| Output: [2, 4, 6, 8] |
| |
| Map applies a function to every element |
+---------------------------------------------+
Mini summary: Map applies a function to every element in a collection.
Definition: Filter is a functional operation that selects only the elements that meet a condition.
Why it is important: Filter helps us get only the data we need.
Simple explanation: Filter is like a sieve that only lets certain things through.
Real-life example: A sieve that only lets small beans through.
School example: A teacher only gives A's to students who scored above 80.
Home example: You only keep the apples that are red.
Nigerian example: A trader only sells goods that are in good condition.
+---------------------------------------------+
| FILTER |
+---------------------------------------------+
| Input: [1, 2, 3, 4, 5, 6] |
| Condition: x > 3 |
| Output: [4, 5, 6] |
| |
| Filter selects only elements that meet |
| the condition |
+---------------------------------------------+
Mini summary: Filter selects only the elements that meet a condition.
Definition: Reduce is a functional operation that combines all elements into a single value.
Why it is important: Reduce helps us summarize data.
Simple explanation: Reduce is like a machine that takes many things and combines them into one.
Real-life example: A machine that takes apples and makes one big apple pie.
School example: A teacher adds all test scores together to get a total.
Home example: You combine all ingredients to make one dish.
Nigerian example: A trader adds all sales to get total revenue.
+---------------------------------------------+
| REDUCE |
+---------------------------------------------+
| Input: [1, 2, 3, 4] |
| Operation: x + y |
| Output: 10 |
| |
| Reduce combines all elements into one |
+---------------------------------------------+
Mini summary: Reduce combines all elements into a single value.
Step 1: Driver creates SparkContext
|
V
Step 2: Driver connects to Cluster Manager
|
V
Step 3: Cluster Manager allocates Executors
|
V
Step 4: Driver sends tasks to Executors
|
V
Step 5: Executors process data in memory
|
V
Step 6: Executors send results to Driver
|
V
Step 7: Driver combines results
|
V
Step 8: Final answer
+---------------------------------------------+
| SPARK ARCHITECTURE |
+---------------------------------------------+
| +---------------------------------------+ |
| | DRIVER | |
| | (Boss – plans the whole job) | |
| +---------------------------------------+ |
| | | | |
| V V V |
| +---------+ +---------+ +---------+ |
| |Executor | |Executor | |Executor | |
| |(Worker) | |(Worker) | |(Worker) | |
| +---------+ +---------+ +---------+ |
| |
| +---------------------------------------+ |
| | CLUSTER MANAGER | |
| | (Assigns tasks to workers) | |
| +---------------------------------------+ |
+---------------------------------------------+
+---------------------------------------------+
| MAP, FILTER, AND REDUCE |
+---------------------------------------------+
| Input: [1, 2, 3, 4, 5, 6] |
| |
| MAP (x -> x * 2): |
| Output: [2, 4, 6, 8, 10, 12] |
| |
| FILTER (x > 3): |
| Output: [4, 5, 6] |
| |
| REDUCE (x + y): |
| Output: 21 |
+---------------------------------------------+
+---------------------------------------------+
| SPARK COMPONENTS |
+---------------------------------------------+
| +---------------------------------------+ |
| | Spark Core (Foundation) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | Spark SQL (Structured data) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | Spark Streaming (Real-time) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | MLlib (Machine Learning) | |
| +---------------------------------------+ |
| +---------------------------------------+ |
| | GraphX (Graph processing) | |
| +---------------------------------------+ |
+---------------------------------------------+
| Feature | Spark | MapReduce |
|---|---|---|
| Speed | Fast (in-memory) | Slow (disk-based) |
| Processing | In-memory | Disk-based |
| Use case | Batch, streaming, ML | Batch only |
| Ease of use | Easier (high-level APIs) | Harder (low-level) |
| Iterative jobs | Excellent | Poor |
| Languages | Python, Scala, Java, R | Java, C++ |
| Operation | What it does | Example |
|---|---|---|
| Map | Applies a function to every element | [1,2,3] -> [2,4,6] |
| Filter | Selects elements that meet a condition | [1,2,3,4] -> [3,4] |
| Reduce | Combines all elements into one | [1,2,3,4] -> 10 |
In this module, we learned about Apache Spark and functional programming.
We started with a story about MapReduce Delivery and Spark Express to understand the speed difference between MapReduce and Spark.
We learned that Apache Spark is a fast, in-memory data processing engine that is much faster than MapReduce.
We explored Spark's architecture and learned about the driver (the boss), executors (the workers), and the cluster manager (the assigner).
We compared Spark to MapReduce and saw that Spark is faster, more versatile, and easier to use.
We learned about the components of Spark: Spark Core, Spark SQL, Spark Streaming, MLlib, and GraphX.
We explored the languages supported by Spark: Python (PySpark), Scala, Java, and R.
We learned about functional programming and how Spark uses it for distributed processing. We explored three common functional operations: Map (apply a function), Filter (select elements), and Reduce (combine elements).
We learned about the Spark shell and how it is the best way to start learning Spark.
We saw examples from Nigeria, including MTN, Jumia, and Flutterwave using Spark.
Remember: Spark is the hero of Big Data processing. It is fast, powerful, and used by thousands of companies around the world.
Match the term on the left with the correct definition on the right.
| Term | Definition |
|---|---|
| 1. Apache Spark | A. The boss that plans the whole job |
| 2. Driver | B. Workers that run tasks |
| 3. Executor | C. Fast, in-memory data processing engine |
| 4. Map | D. Applies a function to every element |
| 5. Filter | E. Selects elements that meet a condition |
| 6. Reduce | F. Combines all elements into one value |
| 7. PySpark | G. Python API for Spark |
Answers: 1-C, 2-A, 3-B, 4-D, 5-E, 6-F, 7-G
Activity: "Design a Spark Solution"
Instructions:
Activity: "My Spark Example"
Instructions:
Title: "Spark for a Nigerian Healthcare System"
Instructions:
Assignment: "Spark Shell Practice"
Instructions:
Title: "Spark vs. MapReduce – The Challenge"
Instructions:
1-C, 2-A, 3-B, 4-D, 5-E, 6-F, 7-G
In the next module, we will learn about Spark RDDs and Data Processing Fundamentals.
We will explore:
To prepare, think about these questions:
We will continue our journey into Big Data by exploring the core data structure of Spark: RDDs. Get ready for an exciting adventure!
End of Module 4
Well done! You have completed the fourth module of Big Data Engineering with Spark.
Module Introduction
Hello, cloud explorer! 🌤️ In Module 4, we learned how Spark helps us process big data fast. But where do we run Spark? We run it on computers. What if we don't have big computers at home? That's where the cloud comes in!
The cloud is like a giant playground of computers that belong to someone else, but we can rent them whenever we want. It's like borrowing your friend's super-powerful gaming computer to play a big game, but you only pay for the time you use it.
In this module, we will learn how to run Spark on the cloud. We will use services like AWS, Azure, and Google Cloud. We will also learn how to set up clusters, manage resources, and make our Spark jobs run even faster. Let's fly into the cloud! 🚀
Learning Objectives
By the end of this module, you will be able to:
Warm‑up Story – The School Playground
Imagine your school has a small playground with just one slide and one swing. When only a few children play, it's fine. But one day, the whole school wants to play at the same time! The playground is too small.
The principal says, "Don't worry! We can use the big playground at the community centre. It has 10 slides, 20 swings, and lots of space. We can rent it for just one day, and we only pay for that day."
That is exactly what cloud computing is! Instead of buying our own big computers (which are expensive and we may not need them every day), we rent them from cloud providers. We pay only for what we use, and we can get as many computers as we need.
Your small computer (at home)
|
| not enough power
V
Cloud playground (big computers)
+-------------------------------+
| Computer 1 Computer 2 |
| Computer 3 Computer 4 |
| ... (hundreds of computers) |
+-------------------------------+
You rent them for a few hours.
🌟 Mini summary: The cloud is like a playground of computers that we can rent when we need them. We pay only for what we use.
Definition: Cloud computing means using computers, storage, and services over the internet, instead of owning them yourself.
Why it is important: The cloud gives us access to powerful computers without buying them. We can scale up or down easily.
Simple explanation: It's like using a library instead of buying all the books yourself. You borrow what you need and return it when done.
Real-life example: Netflix uses the cloud to stream movies to millions of people.
School example: Your school uses Google Classroom – that's cloud computing!
Home example: You use iCloud or Google Photos to store your pictures.
Nigerian example: Many Nigerian startups use the cloud to run their apps without buying expensive servers.
Traditional: Buy computer → Install software → Run Cloud: Rent computer over internet → Run → Stop paying
📌 Mini summary: Cloud computing means renting computers and services over the internet instead of buying them.
Definition: The three biggest cloud companies are Amazon Web Services (AWS), Microsoft Azure, and Google Cloud Platform (GCP).
Why it is important: These companies have the biggest and most reliable clouds. Most people use one of them.
Simple explanation: They are like the three biggest supermarkets for computers. You can go to any one to buy (rent) computer power.
Real-life example: Spotify uses Google Cloud. Airbnb uses AWS. Many banks use Azure.
School example: Your school might use Microsoft Office 365 (Azure) or Google Workspace (GCP).
Home example: If you use Amazon Prime Video, that's AWS.
Nigerian example: Flutterwave uses AWS for its payment processing.
+----------+ +----------+ +----------+ | AWS | | Azure | | GCP | | (Amazon) | | (Microsoft) | (Google) | +----------+ +----------+ +----------+ All three offer Spark in the cloud!
📌 Mini summary: AWS, Azure, and GCP are the three big cloud providers. They all offer Spark.
Definition: Running Spark in the cloud means you run your Spark jobs on rented computers that are managed by a cloud provider.
Why it is important: You don't need to buy expensive hardware. You can get thousands of computers for a short time.
Simple explanation: It's like having a magic button that gives you more computers when you need them.
Real-life example: A company runs a big sale and needs extra computer power for one day. They rent it from the cloud.
School example: During exam grading, your school rents cloud computers to process all the results quickly.
Home example: If your family wants to store a million photos, they use cloud storage.
Nigerian example: During elections, cloud computers can count voting data from all over the country.
Reasons to use cloud for Spark: 1. No need to buy hardware 💰 2. Scale up or down easily 📈 3. Pay only for what you use ⏱️ 4. Use the latest machines 🆕 5. Focus on your code, not on setup 🧑💻
📌 Mini summary: The cloud makes Spark easy and affordable because you rent computers only when you need them.
Definition: AWS is Amazon's cloud platform. It is the most popular cloud provider in the world.
Why it is important: AWS has many services for big data, including EMR (Elastic MapReduce) which runs Spark.
Simple explanation: AWS is like a huge toy store with many different toys (services). One of the toys is Spark.
Real-life example: Netflix, Airbnb, and many others use AWS.
School example: If your school uses Amazon, they might use AWS for their website.
Home example: Amazon Prime Video runs on AWS.
Nigerian example: Paystack (owned by Stripe) uses AWS for payment processing.
AWS Services for Big Data: +------------------+ | EMR (Spark) | ← we will use this! | S3 (storage) | ← store our data | EC2 (computers) | ← virtual machines | Lambda (serverless) | +------------------+
📌 Mini summary: AWS is Amazon's cloud. We use EMR to run Spark and S3 to store data.
Definition: S3 is a service in AWS that stores files in the cloud. Think of it as a very big online folder.
Why it is important: S3 can store any amount of data. Spark can read data directly from S3.
Simple explanation: It's like Google Drive but made for big data. You can put any file there, and Spark can read it.
Real-life example: Companies store all their customer data in S3.
School example: Your teacher stores all assignments in a shared folder (like S3).
Home example: You save your homework in Google Drive – S3 is similar but more powerful.
Nigerian example: A bank stores all transaction logs in S3.
Data Flow: Your computer → upload to S3 → Spark reads from S3 → results saved to S3
📌 Mini summary: S3 is cloud storage for files. Spark can read and write data to S3.
Definition: EMR is an AWS service that makes it easy to run Spark and other big data tools in the cloud.
Why it is important: EMR sets up all the Spark clusters for you automatically. You just tell it what to do.
Simple explanation: EMR is like a magical machine that creates Spark clusters with a single click. No need to set up anything yourself.
Real-life example: A data scientist uses EMR to run a Spark job on 100 computers in 5 minutes.
School example: The principal clicks a button and suddenly 20 computers are ready for the students.
Home example: Like ordering a pizza – you click, and it arrives ready.
Nigerian example: A startup uses EMR to process customer data every night.
How EMR works: 1. You tell EMR what you need (Spark, number of computers, etc.) 2. EMR creates the cluster automatically 3. You run your Spark job 4. EMR shuts everything down when done (optional)
📌 Mini summary: EMR is a service that creates and manages Spark clusters in the cloud for you.
Definition: Azure has a service called Synapse that can run Spark. Databricks is a company that makes Spark easy on all clouds.
Why it is important: You don't have to use AWS – you can use Azure or Databricks too.
Simple explanation: It's like having different brands of cars – all can take you to the same place (Spark), but each has its own style.
Real-life example: Many companies use Databricks because it has a nice web interface.
School example: Your school might use different software for different subjects.
Home example: You can use different apps to watch movies.
Nigerian example: A Nigerian company might choose Azure because they already use Microsoft products.
Cloud Spark Options:
+----------+ +----------+ +----------+
| AWS EMR | | Azure | | Google |
| | | Synapse | | Dataproc |
+----------+ +----------+ +----------+
| | |
+---------------+---------------+
|
V
+-------------+
| Databricks | (works on all clouds)
+-------------+
📌 Mini summary: You can run Spark on AWS, Azure, or Google Cloud. Databricks is a popular tool that works on all of them.
Definition: A cluster is a group of computers that work together. Setting up a cluster means creating those computers.
Why it is important: Without a cluster, you can't run Spark in parallel. You need many workers.
Simple explanation: It's like forming a team. You need to gather the players before you can play the game.
Real-life example: Before a big project, the manager forms a team of workers.
School example: The teacher divides the class into groups before the activity.
Home example: Before cleaning the house, you ask everyone to help.
Nigerian example: Before the village harvest, the chief calls all the farmers together.
Steps to create an EMR cluster: 1. Go to AWS Management Console 2. Click on EMR 3. Click "Create Cluster" 4. Choose Spark as the application 5. Choose number of nodes (computers) 6. Click "Create" – wait 5 minutes 7. Your cluster is ready! 🎉
📌 Mini summary: Creating a cluster on EMR is easy – you click a few buttons and wait a few minutes.
Definition: Reading data from S3 means your Spark code can load files directly from S3 storage.
Why it is important: Your data is in S3, so Spark needs to read it from there.
Simple explanation: It's like reading a book from the library – you need to go to the library first.
Real-life example: A company stores all sales data in S3. Spark reads it to analyse sales.
School example: The teacher stores all homework in a cloud folder. Students read from there.
Home example: You read recipes from an online cookbook.
Nigerian example: A farmer stores crop data in S3. Spark reads it to predict harvest.
Spark code to read from S3:
df = spark.read.csv("s3://my-bucket/sales_data.csv")
or
df = spark.read.parquet("s3://my-bucket/customers/")
S3 path format: s3://bucket-name/folder/file
📌 Mini summary: Spark can read files from S3 using paths that start with "s3://".
Definition: Writing results to S3 means saving your Spark output to S3 storage.
Why it is important: You want to keep your results for later. S3 is a safe place to store them.
Simple explanation: After you finish your drawing, you put it in your folder so you don't lose it.
Real-life example: After running a report, the company saves it to S3 for other teams to use.
School example: After the class project, the teacher saves it to the school's cloud folder.
Home example: After you finish your homework, you save it to your computer.
Nigerian example: After counting votes, the result is saved to S3 for everyone to see.
Spark code to write to S3:
result_df.write.csv("s3://my-bucket/output/results.csv")
or
result_df.write.parquet("s3://my-bucket/output/")
You can also save as JSON, ORC, etc.
📌 Mini summary: Spark can save results to S3 using the write() method.
Definition: Scaling means using more or fewer computers to do the work.
Why it is important: Sometimes you need more power (more computers) to finish fast. Other times you want fewer computers to save money.
Simple explanation: It's like adding more people to help you move a heavy sofa. More people = faster, but you have to pay them.
Real-life example: During Black Friday, an online store uses 10 times more computers to handle all the orders.
School example: For a big exam, the school uses more teachers to mark papers.
Home example: When cleaning the house, more family members help to finish quickly.
Nigerian example: During elections, INEC uses more computers to count votes quickly.
Scaling options: +------------------+------------------+ | Vertical Scaling | Horizontal Scaling | | (make one bigger) | (add more workers) | +------------------+------------------+ | Use a bigger | Use 10 computers | | computer | instead of 1 | +------------------+------------------+ In the cloud, horizontal scaling is easier!
📌 Mini summary: Scaling means adding or removing computers to match your needs.
Definition: Monitoring means watching your Spark job while it runs to see if everything is going well.
Why it is important: You want to know if your job is making progress, using enough memory, or if it has errors.
Simple explanation: It's like watching the cooker to see if your food is cooking properly.
Real-life example: A driver uses a speedometer to monitor the car's speed.
School example: The teacher walks around the class to monitor students' work.
Home example: Mum checks the oven to see if the cake is ready.
Nigerian example: A bus conductor monitors how many passengers are on the bus.
Monitoring tools: 1. EMR Console (shows cluster status) 2. Spark UI (shows job progress) 3. CloudWatch (shows logs and metrics) 4. Ganglia (shows resource usage) You can see: - How much memory is used - How many tasks are running - If any tasks failed
📌 Mini summary: Monitoring helps you see if your Spark job is running correctly.
Definition: Cost management means using the cloud in a way that doesn't waste money.
Why it is important: Cloud services cost money per hour. If you leave computers running, you pay even if you're not using them.
Simple explanation: It's like leaving all the lights on in your house – you pay for electricity you don't need.
Real-life example: A company shuts down its cloud computers at night to save money.
School example: The teacher turns off the projector when not in use.
Home example: You switch off the TV when no one is watching.
Nigerian example: A small business uses cloud computers only during business hours to save money.
Ways to save money: 1. ✅ Shut down clusters when not in use 2. ✅ Use spot instances (cheaper) 3. ✅ Choose smaller machines if you don't need big ones 4. ✅ Use auto-scaling (add computers only when needed) 5. ❌ Don't leave clusters running overnight (unless you need to)
📌 Mini summary: To save money, turn off cloud computers when you don't need them.
Definition: Serverless means you don't manage servers at all. You just submit your Spark job and it runs somewhere.
Why it is important: It makes Spark even easier. You don't need to create clusters or manage anything.
Simple explanation: It's like using a vending machine – you put in your request (code) and get the result (output). You don't care about how the machine works inside.
Real-life example: AWS Glue is a serverless Spark service. You just write your ETL jobs and they run.
School example: The canteen gives you food – you don't need to know how it's cooked.
Home example: You use a microwave – you don't know how it works inside.
Nigerian example: A payment app like Opay processes transactions – you don't see the servers.
Serverless vs Traditional: Traditional: You manage cluster, software, updates Serverless: You just submit code – the cloud handles everything +------------------+------------------+ | Traditional | Serverless | | (EMR) | (Glue, Databricks) | +------------------+------------------+
📌 Mini summary: Serverless Spark means you don't manage any servers – just submit your code.
Here is how a typical cloud Spark workflow looks:
1. Prepare your data → upload to S3 2. Write your Spark code (Python/Scala) 3. Create an EMR cluster (or use serverless) 4. Submit your job to the cluster 5. Monitor progress via Spark UI 6. Results are saved to S3 7. Shut down the cluster (to save money!) 8. Download results from S3 (if needed)
Let's see a simple example in Python using PySpark on EMR:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("Cloud Word Count").getOrCreate()
# Read from S3
df = spark.read.text("s3://my-bucket/story.txt")
# Process
words = df.rdd.flatMap(lambda line: line[0].split(" "))
word_counts = words.map(lambda w: (w, 1)).reduceByKey(lambda a,b: a+b)
# Save to S3
word_counts.toDF(["word", "count"]).write.csv("s3://my-bucket/output/")
spark.stop()
📌 Mini summary: Cloud Spark workflow: upload data → write code → create cluster → run → save results → shut down.
How to run a Spark job on EMR:
Step 1: AWS Console → EMR Step 2: Create Cluster Step 3: Add Step (your Spark code) Step 4: Run Job Step 5: Check Output Step 6: Terminate Cluster 💰
Cloud Architecture (simple):
+-------------+
| Your Laptop|
+-------------+
|
| (submit job)
V
+---------------------------------+
| AWS CLOUD |
| +---------------------------+ |
| | EMR Cluster | |
| | +------+ +------+ | |
| | |Master| |Worker| | |
| | +------+ +------+ | |
| | |Worker| |Worker| | |
| | +------+ +------+ | |
| +---------------------------+ |
| +---------------------------+ |
| | S3 Storage | |
| | (data in, data out) | |
| +---------------------------+ |
+---------------------------------+
|
V
Results saved to S3
Scaling Up vs Scaling Out:
Scaling Up (Vertical): [Small computer] → [Bigger computer] (one computer, more powerful) Scaling Out (Horizontal): [1 computer] → [10 computers] (more computers, each does a small part)
Cost vs Time Trade‑off:
+------------------+------------------+ | Number of Nodes | Time to Finish | +------------------+------------------+ | 1 node | 10 hours | | 5 nodes | 2 hours | | 10 nodes | 1 hour | | 20 nodes | 30 minutes | +------------------+------------------+ More nodes = faster, but more expensive.
On‑Premise vs Cloud:
| Feature | On‑Premise | Cloud |
|---|---|---|
| Cost | High upfront | Pay as you go |
| Scaling | Hard and slow | Easy and fast |
| Maintenance | You do it | Provider does it |
| Setup time | Weeks | Minutes |
| Accessibility | On‑site only | From anywhere |
AWS vs Azure vs GCP:
| Feature | AWS | Azure | GCP |
|---|---|---|---|
| Spark service | EMR | Synapse | Dataproc |
| Storage | S3 | Blob | Cloud Storage |
| Popularity | #1 | #2 | #3 |
| Free tier | Yes | Yes | Yes |
| Serverless Spark | Glue | Synapse | Dataproc Serverless |
EMR Instance Types:
| Type | Use case | Cost |
|---|---|---|
| m5.xlarge | General purpose | Medium |
| r5.xlarge | Memory‑intensive | Higher |
| c5.xlarge | Compute‑intensive | Medium |
| spot instances | Non‑critical jobs | Very low |
Congratulations! 🎉 You have completed Module 5 – Spark in the Cloud. Here is what we learned:
You are now ready for Module 6 – Real‑World Spark Projects! 🚀
Match the term with its description:
| Term | Description |
|---|---|
| 1. AWS | A. Amazon's cloud |
| 2. S3 | B. Cloud storage |
| 3. EMR | C. Spark in the cloud |
| 4. Glue | D. Serverless Spark |
| 5. EC2 | E. Virtual machines |
Answers: 1-A, 2-B, 3-C, 4-D, 5-E
Title: Cloud Provider Comparison
Instructions: In groups of 4, research one cloud provider (AWS, Azure, GCP). Prepare a 5‑minute presentation about:
Title: My First S3 Upload
Instructions: If you have AWS access, create an S3 bucket, upload a small CSV file, and read it with Spark (using EMR or Glue). Write down the steps you followed.
Title: Cloud‑Based Sales Analytics
Description: You have sales data for a Nigerian supermarket stored in S3. Use Spark on EMR to:
Deliverable: Spark script and a brief report on how you set up the cluster.
Title: Run Word Count on EMR
Instructions:
Title: Optimise Cloud Costs
Problem: Your company runs a Spark job every day that takes 4 hours on a 10‑node cluster. The cluster costs ₦10,000 per hour. You want to reduce costs by 50%. Suggest two strategies and explain how they would work.
Fill‑in‑the‑Blank Answers: 1. internet, 2. Amazon, 3. storage, 4. MapReduce, 5. servers
True/False Answers: 1. False, 2. False, 3. True, 4. False, 5. True
Multiple Choice Answers: 1-B, 2-C, 3-B, 4-B, 5-A, 6-B, 7-A, 8-B, 9-B, 10-C, 11-B, 12-B, 13-A, 14-B, 15-A
Matching Answers: 1-A, 2-B, 3-C, 4-D, 5-E
In Module 6, we will apply everything we have learned to real‑world projects. We will build an end‑to‑end data pipeline, work with real datasets, and solve actual business problems.
What to bring:
See you in Module 6 – let's build something great! 🏗️
Module Introduction
Welcome, young data engineer! 🎉 You have made it to the final module of our Spark journey. In Modules 1 to 5, we learned what Spark is, how it works, and how to run it in the cloud. Now it's time to put everything together and build real-world projects!
Imagine you are a chef who has learned all the cooking techniques. Now it's time to cook a full meal! In this module, we will build complete projects from start to finish. We will solve problems that real companies face every day.
By the end of this module, you will be able to build your own Spark projects and show them to your friends, teachers, and maybe even future employers. Let's build something amazing! 🚀
Learning Objectives
By the end of this module, you will be able to:
Warm‑up Story – The School Sports Day
Your school is organising a big sports day with many events: running, jumping, and throwing. The PE teacher needs to know which house (Red, Blue, Green, Yellow) is winning overall. There are hundreds of students, and each one participates in multiple events.
You decide to help the teacher. You plan everything:
This is exactly how real-world Spark projects work! You follow the same steps to solve business problems.
Define Problem → Collect Data → Clean Data → Compute Results → Present Findings
(1) (2) (3) (4) (5)
🌟 Mini summary: Every data project follows 5 steps: Define, Collect, Clean, Compute, Present.
Definition: A data project is a project that uses data to solve a problem or answer a question.
Why it is important: Every Spark project follows the same pattern. Knowing the steps helps you stay organised.
Simple explanation: It's like following a recipe when cooking – you do things in a certain order to get a good result.
Real-life example: A supermarket wants to know which products sell best. They follow these 5 steps.
School example: You want to find out which subject your class enjoys most. You survey everyone and follow the steps.
Home example: Your family wants to know which meal is everyone's favourite. You collect votes and count them.
Nigerian example: A bank wants to know which ATM machine is most used in Lagos.
5 Steps of a Data Project: +------------------+-----------------------------------+ | Step 1: Define | What question are we answering? | | Step 2: Ingest | Where is the data? | | Step 3: Transform| Clean and prepare the data. | | Step 4: Analyse | Compute answers and insights. | | Step 5: Visualise| Show results clearly. | +------------------+-----------------------------------+
📌 Mini summary: Every data project has 5 steps: Define, Ingest, Transform, Analyse, Visualise.
Definition: Defining the problem means understanding what question you need to answer.
Why it is important: If you don't know what you're looking for, you won't find it!
Simple explanation: Before you start building, you need to know: "What do I want to find out?"
Real-life example: A mobile network wants to know: "Which areas have the worst call drops?"
School example: "Which student read the most books this term?"
Home example: "Which day of the week do we spend the most money?"
Nigerian example: "Which state has the highest number of farmers registered?"
Define the Problem: +----------------------------------+ | Ask: What is the business question? | | Ask: Why is this important? | | Ask: Who will use the answer? | | Write it down clearly. | +----------------------------------+
📌 Mini summary: Defining the problem is the most important step. Know what you're trying to find out.
Definition: Ingesting data means bringing data into your system from where it is stored.
Why it is important: You can't analyse data you don't have!
Simple explanation: This is like going to the market to buy the ingredients you need for your recipe.
Real-life example: A company downloads sales data from their database into Spark.
School example: The teacher collects all the test scores from each class.
Home example: You collect all the receipts from the past month.
Nigerian example: INEC collects voting results from all polling stations.
Ingest Data: Data sources: - CSV files (spreadsheets) - JSON files (web data) - Databases (SQL) - S3 (cloud storage) - HDFS (Hadoop storage) Spark can read from all of these!
📌 Mini summary: Ingest means bringing data into Spark from wherever it is stored.
Definition: Transforming data means cleaning, filtering, and changing it so it's ready for analysis.
Why it is important: Real data is messy. It has missing values, errors, and inconsistencies.
Simple explanation: This is like washing and cutting vegetables before you cook them.
Real-life example: Removing duplicate customer records from a database.
School example: Fixing students' names that were spelled differently in different records.
Home example: Organising all your toys by colour and size.
Nigerian example: Correcting different spellings of "Lagos" (e.g., "Lagos", "Legos").
Transform Data (Cleaning): 1. Remove duplicates 2. Handle missing values (fill with average or delete) 3. Fix inconsistent formatting (e.g., date formats) 4. Filter out irrelevant data 5. Convert data types (e.g., string to number)
📌 Mini summary: Transform means cleaning and preparing data so it's ready for analysis.
Definition: Analysing data means using Spark to compute answers to your questions.
Why it is important: This is where you get the actual results!
Simple explanation: This is like cooking the meal – you put everything together and get the final dish.
Real-life example: Calculating total sales by region.
School example: Finding the average score for each subject.
Home example: Adding up all your monthly expenses.
Nigerian example: Counting votes per candidate in an election.
Analyse Data (Compute): - Aggregations (sum, count, average) - Joins (combining different datasets) - Filtering (selecting specific data) - Sorting (ordering results) - Grouping (data by category)
📌 Mini summary: Analyse means using Spark to compute answers from your clean data.
Definition: Visualising means showing your results in a way that is easy to understand, like charts or dashboards.
Why it is important: People need to see the results clearly. Numbers alone can be confusing.
Simple explanation: This is like putting your food on a nice plate so it looks good to eat.
Real-life example: Creating a bar chart showing sales by month.
School example: Drawing a graph of class test scores.
Home example: Making a chart of your weekly chores progress.
Nigerian example: Showing election results on a map of Nigeria.
Visualise and Present: - Bar charts (for comparisons) - Line charts (for trends over time) - Pie charts (for proportions) - Tables (for detailed data) - Maps (for geographic data)
📌 Mini summary: Visualise means showing your results in charts and graphs that are easy to understand.
Let's build our first real project: a sales analytics dashboard for a store.
Problem: A supermarket wants to know which products sell best, which days are busiest, and which customers spend the most.
Data: Sales records with columns: date, product, price, quantity, customer_id.
Steps:
Step 1 – Define:
Question: What are the top 10 products by sales?
Step 2 – Ingest:
sales_df = spark.read.csv("s3://my-bucket/sales.csv", header=True)
Step 3 – Transform:
sales_df = sales_df.filter(sales_df.price > 0) # remove bad data
sales_df = sales_df.dropDuplicates() # remove duplicates
Step 4 – Analyse:
sales_df.createOrReplaceTempView("sales")
top_products = spark.sql("""
SELECT product, SUM(price * quantity) as total_sales
FROM sales
GROUP BY product
ORDER BY total_sales DESC
LIMIT 10
""")
Step 5 – Visualise:
top_products.show() # print the table
# (we would use a plotting library to make charts)
📌 Mini summary: Sales analytics helps businesses understand what sells best and when.
Let's build a movie recommendation system like Netflix or Amazon Prime.
Problem: Suggest movies that a user might like based on what they've watched before.
Data: User ratings: user_id, movie_id, rating (1-5), timestamp.
Steps:
Step 1 – Define:
Question: What movies should we recommend to user 123?
Step 2 – Ingest:
ratings_df = spark.read.csv("ratings.csv", header=True)
Step 3 – Transform:
ratings_df = ratings_df.filter(ratings_df.rating >= 1) # valid ratings
Step 4 – Analyse:
from pyspark.ml.recommendation import ALS
als = ALS(maxIter=5, regParam=0.01, userCol="user_id",
itemCol="movie_id", ratingCol="rating")
model = als.fit(ratings_df)
recommendations = model.recommendForAllUsers(5)
Step 5 – Visualise:
recommendations.show() # show top 5 movies for each user
How it works: The ALS (Alternating Least Squares) algorithm finds patterns in who likes what movies. Then it predicts what movies a user will enjoy based on patterns from similar users.
User A likes: Movie 1, Movie 2, Movie 3 User B likes: Movie 1, Movie 2, Movie 4 User C likes: Movie 1, Movie 4, Movie 5 → Movie 1 is popular → Users who like Movie 1 also like Movie 4 → Recommend Movie 4 to User A!
📌 Mini summary: Recommendation systems use patterns in data to suggest things people might like.
Let's build a system to analyse server logs and monitor for problems.
Problem: A website wants to know when it has errors, how many visitors it gets, and which pages are popular.
Data: Server logs: timestamp, ip_address, page_url, status_code, response_time.
Steps:
Step 1 – Define:
Question: How many 404 (page not found) errors occur each hour?
Step 2 – Ingest:
logs_df = spark.read.text("s3://my-bucket/logs/")
# Parse the log lines (complex parsing needed)
Step 3 – Transform:
# Parse each log line into columns
logs_df = logs_df.withColumn("timestamp", ...)
logs_df = logs_df.withColumn("status", ...)
logs_df = logs_df.filter(logs_df.status == 404)
Step 4 – Analyse:
logs_df.groupBy("hour").count().orderBy("hour")
Step 5 – Visualise:
# Show a chart of errors by hour
Common log format (Apache):
192.168.1.1 - - [01/Jan/2024:12:00:00] "GET /index.html" 200 1024 This shows: IP, timestamp, request, status code, size
📌 Mini summary: Log analysis helps monitor websites and catch problems before they affect users.
Definition: Clean code means code that is easy to read, understand, and maintain.
Why it is important: Other people (and future you!) need to understand your code. Clean code saves time.
Simple explanation: It's like writing neatly so your teacher can read your homework.
Real-life example: A team of 10 data engineers works on the same code. Clean code helps them collaborate.
School example: You write your answers clearly so the teacher can mark them.
Home example: You label your toy boxes so you know what's inside.
Nigerian example: A team in a Nigerian startup works on the same project – clean code helps everyone.
Tips for Clean Spark Code:
1. Use meaningful variable names (not "x", "y" but "sales_df")
2. Add comments explaining what each part does
3. Break long code into smaller functions
4. Use consistent formatting (spaces, indentation)
5. Follow a style guide (like PEP 8 for Python)
Bad:
df = spark.read.csv("file.csv")
d = df.filter(df.a > 0)
d.show()
Good:
sales_df = spark.read.csv("sales_data.csv", header=True)
filtered_sales = sales_df.filter(sales_df.price > 0)
filtered_sales.show()
📌 Mini summary: Clean code is easy to read and understand. It helps you and others work better.
Definition: Testing means checking if your code works correctly. Debugging means finding and fixing errors.
Why it is important: Bugs (errors) can cause wrong results, which can lead to bad decisions.
Simple explanation: It's like checking your homework before submitting it.
Real-life example: A bank tests its fraud detection system to make sure it catches fraud.
School example: You review your answers before handing in the test.
Home example: You taste the food before serving it to check if it's good.
Nigerian example: A telecom company tests its network monitoring system.
Testing and Debugging Tips: 1. Test with a small sample of data first 2. Use assert statements to check results 3. Print intermediate results (for debugging) 4. Check the Spark UI for errors 5. Use try/except to handle errors gracefully 6. Write unit tests for your functions
📌 Mini summary: Testing and debugging help you find and fix errors before they cause problems.
Definition: Documentation means writing down what your code does and how to use it.
Why it is important: Documentation helps others (and future you) understand the project.
Simple explanation: It's like writing instructions for a game so others can play it too.
Real-life example: A software company has a "user manual" for its product.
School example: You write a summary of a book to help others understand it.
Home example: You write a note to remind your family how to use the new TV.
Nigerian example: A Nigerian tech team documents their code so new members can learn quickly.
What to Document: 1. What is the problem being solved? 2. What data is used? 3. What are the main steps in the code? 4. How do you run the code? 5. What are the expected outputs? 6. Any special requirements or dependencies? Example: # This program calculates total sales by region. # Input: sales_data.csv with columns: region, product, price, quantity # Output: CSV file with region and total sales # Run: spark-submit sales_analysis.py
📌 Mini summary: Documentation helps others understand and use your code.
Definition: Presenting means sharing your findings in a clear and engaging way.
Why it is important: If you can't explain your results, people won't understand the value of your work.
Simple explanation: It's like telling a story – you want people to understand and remember your message.
Real-life example: A data scientist presents their findings to the CEO.
School example: You present your project to the class.
Home example: You explain to your family why you want a pet.
Nigerian example: A startup pitches their data product to investors.
Tips for Presenting Results: 1. Start with the question/problem 2. Show the data (where it came from) 3. Explain the steps you took 4. Show the results clearly (charts, tables) 5. Tell what the results mean 6. Make recommendations 7. Keep it simple and avoid jargon 8. Practice your presentation!
📌 Mini summary: Presenting your results means telling a story about what you found and why it matters.
Definition: The project lifecycle is the journey from the initial idea to the final product.
Why it is important: Understanding the lifecycle helps you plan and manage your projects.
Simple explanation: It's like the journey of building a house – from drawing plans to moving in.
Real-life example: A company goes through many stages to launch a new product.
School example: Your school project goes from idea to research to final presentation.
Home example: Planning a party – from idea to invitations to food to the actual party.
Nigerian example: A Nigerian startup builds an app from idea to launch.
Project Lifecycle: +------------------+-----------------------------+ | 1. Idea | What problem to solve? | | 2. Planning | What data and tools? | | 3. Development | Write the code | | 4. Testing | Is it correct? | | 5. Deployment | Run it in production | | 6. Monitoring | Is it working? | | 7. Maintenance | Fix and improve | +------------------+-----------------------------+
📌 Mini summary: The project lifecycle takes an idea through planning, building, testing, and running.
Let's build a complete project from start to finish: Nigerian Food Sales Analysis
Problem: A restaurant chain wants to know which foods are most popular in different states.
Data: Sales records: restaurant_id, state, food_item, quantity, price, date.
Complete Code:
# Step 1: Import Spark
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("Nigerian Food Sales").getOrCreate()
# Step 2: Ingest Data
sales_df = spark.read.csv("s3://my-bucket/food_sales.csv", header=True, inferSchema=True)
# Step 3: Transform Data
# Remove rows with missing values
sales_df = sales_df.dropna()
# Add a total_sales column
sales_df = sales_df.withColumn("total_sales", sales_df.quantity * sales_df.price)
# Step 4: Analyse Data
# Find top foods by state
sales_df.createOrReplaceTempView("sales")
top_foods = spark.sql("""
SELECT state, food_item, SUM(total_sales) as total_sales
FROM sales
GROUP BY state, food_item
ORDER BY state, total_sales DESC
""")
# Get overall top 10 foods
overall_top = spark.sql("""
SELECT food_item, SUM(total_sales) as total_sales
FROM sales
GROUP BY food_item
ORDER BY total_sales DESC
LIMIT 10
""")
# Step 5: Visualise and Present
print("Overall Top 10 Foods:")
overall_top.show()
# Save results
top_foods.write.csv("s3://my-bucket/top_foods_by_state/")
overall_top.write.csv("s3://my-bucket/overall_top_foods/")
print("Analysis complete! 🎉")
spark.stop()
Sample Output: +-------+------------+------------+ | state | food_item | total_sales| +-------+------------+------------+ | Lagos | Jollof Rice| 1,245,000 | | Lagos | Egusi Soup | 987,000 | | Abuja | Suya | 654,000 | | Port H| Pepper Soup| 432,000 | +-------+------------+------------+
📌 Mini summary: A complete project takes you from problem definition to final results.
How to build a Spark project from scratch:
Project Building Steps: 1. Understand the problem 2. Find the data 3. Explore the data 4. Ingest with Spark 5. Transform (clean) 6. Analyse 7. Save results 8. Visualise 9. Present findings 10. Get feedback and improve
Project Lifecycle Flow:
+---------+ +---------+ +---------+
| Define | --> | Ingest | --> |Transform|
+---------+ +---------+ +---------+
| |
| |
V V
+---------+ +---------+ +---------+
|Visualise| <-- | Analyse | <-- |(repeat) |
+---------+ +---------+ +---------+
End-to-End Pipeline:
Raw Data → Clean Data → Analysed Data → Results → Dashboard (messy) (spark code) (spark code) (output) (charts)
Recommendation System Flow:
User Ratings → Train Model → Predict Scores → Recommend (history) (ALS algorithm) (for each user) (top items)
Log Analysis Pipeline:
Server Logs → Parse Logs → Filter Errors → Group by Hour → Alert (raw text) (extract data) (status=404) (count per hour) (if > threshold)
Project Types Comparison:
| Project Type | Goal | Output | Users |
|---|---|---|---|
| Sales Analytics | Understand sales trends | Charts, dashboards | Management |
| Recommendation | Suggest products/content | Recommendations | Customers |
| Log Analysis | Monitor system health | Alerts, reports | IT, DevOps |
| Fraud Detection | Find suspicious activity | Alerts, flags | Security, Risk |
Data Source Comparison:
| Data Source | Example | Best For | Spark Read |
|---|---|---|---|
| CSV | sales_data.csv | Simple tables | spark.read.csv |
| JSON | events.json | Nested data | spark.read.json |
| Parquet | data.parquet | Large datasets | spark.read.parquet |
| Database | MySQL, PostgreSQL | Live data | spark.read.jdbc |
| S3 | s3://bucket/data/ | Cloud storage | spark.read.csv("s3://...") |
Congratulations! 🎉 You have completed Module 6 – Real‑World Spark Projects. This is the final module of our Spark journey. Let's summarise everything we learned:
You now have all the skills to build your own Spark projects! Keep learning, keep building, and keep asking questions. The world of big data is waiting for you! 🌍
Match the term with its description:
| Term | Description |
|---|---|
| 1. Define | A. Clean and prepare data |
| 2. Ingest | B. Show results with charts |
| 3. Transform | C. Know what problem to solve |
| 4. Analyse | D. Bring data into Spark |
| 5. Visualise | E. Compute answers from data |
Answers: 1-C, 2-D, 3-A, 4-E, 5-B
Title: Design a Data Project
Instructions: In groups of 4, design a complete data project for a Nigerian business (e.g., a restaurant, a transport company, a school). Define the problem, describe the data, and outline the steps you would take. Present your design to the class.
Title: My Own Data Project
Instructions: Think of a problem you want to solve in your school or community. Write a one‑page proposal that includes:
Title: Nigerian Food Sales Dashboard
Description: Build a complete sales analytics project using Spark. You are given sales data for a Nigerian restaurant chain. Your tasks:
Deliverable: Spark code, output files, and a presentation.
Title: Build a Movie Recommendation System
Instructions:
Title: Real-Time Log Monitoring
Problem: Build a system that monitors server logs in real time and alerts if more than 100 errors occur in 5 minutes.
Hint: Use Spark Streaming to process logs as they arrive. Use a sliding window of 5 minutes. If the error count exceeds 100, print an alert.
Fill‑in‑the‑Blank Answers: 1. define, 2. ingest, 3. transform, 4. visualise, 5. clean
True/False Answers: 1. True, 2. False, 3. False, 4. True, 5. False
Multiple Choice Answers: 1-B, 2-B, 3-A, 4-B, 5-B, 6-A, 7-B, 8-B, 9-B, 10-B, 11-A, 12-B, 13-A, 14-B, 15-B
Matching Answers: 1-C, 2-D, 3-A, 4-E, 5-B
Congratulations on completing all 6 modules of "Big Data Engineering with Spark"! 🎉
You have learned:
What you can do next:
Remember: The best way to learn is by doing. Keep building, keep experimenting, and keep asking questions. You are now a data engineer – go and change the world with data! 🌍🚀
Module Introduction
Hello, real-time data explorer! ⚡ So far in our Spark journey, we have been working with data that is already stored in files or databases. That's called batch processing – we process data in chunks, like reading a whole book at once.
But what if data keeps coming all the time, like water flowing from a tap? What if we want to analyse it as it arrives – second by second? That's called real-time processing or streaming.
Imagine you are watching a football match and you want to count how many goals are scored as they happen. You don't wait until the match is over. You count them live! That's exactly what Spark Streaming does – it processes data live, as it arrives.
In this module, we will learn about Spark Streaming, how it works, and how to build real-time applications. Let's dive into the stream! 🌊
Learning Objectives
By the end of this module, you will be able to:
Warm‑up Story – The Busy Bakery
Remember Mama Chidi's bakery from Module 4? Her bakery is now very popular! Every minute, customers buy bread, cakes, and pastries. Mama Chidi wants to know right now how many items are being sold, which items are most popular, and if they are running out of stock.
She can't wait until the end of the day to count everything. She needs to know as it happens. So she installs a computer system at the cash register that sends data to Spark every time a customer buys something.
Spark receives this data continuously, like a river of information. It counts the items, updates the dashboard, and even sends an alert if stock is low. This is streaming data processing!
Customer buys bread → Cash register sends data → Spark processes → Dashboard updates (every second) (continuously) (in real-time) (live view)
🌟 Mini summary: Streaming processes data as it arrives, giving us real-time insights.
Definition: Streaming data is data that is generated continuously and arrives in small pieces over time.
Why it is important: Many things in life happen continuously. We need to analyse them as they happen, not wait until later.
Simple explanation: Think of a river. Water keeps flowing. You can't wait for the whole river to pass – you need to analyse the water as it flows.
Real-life example: Tweets on Twitter. Every second, thousands of new tweets are posted.
School example: Students entering the school gate every morning. The gate counts them as they enter.
Home example: Water dripping from a tap. Each drop is data.
Nigerian example: Mobile money transactions on Opay. Every transaction is a piece of streaming data.
Batch Data (static): [all data at once] → process → result Streaming Data (live): [data1, data2, data3, ...] → process continuously → results
📌 Mini summary: Streaming data comes continuously, like water from a tap.
Definition: Batch processing means collecting data over a period and processing it all at once. Stream processing means processing data as it arrives.
Why it is important: Different problems need different approaches. Some need real-time answers; others can wait.
Simple explanation: Batch is like washing all your clothes at the end of the week. Stream is like washing each item as soon as it gets dirty.
Real-life example: A bank processes daily transactions in a batch at night (batch). It also detects fraud instantly as transactions happen (stream).
School example: Grading all exams at the end of the term (batch). Taking attendance each morning (stream).
Home example: Cleaning the whole house on Saturday (batch). Wiping the kitchen counter after each meal (stream).
Nigerian example: Counting votes after the election (batch). Monitoring election results as they are announced (stream).
+-----------------------+-----------------------+ | Batch Processing | Stream Processing | +-----------------------+-----------------------+ | Process at once | Process continuously | | Wait for all data | Process as it arrives | | Good for large files | Good for live data | | Example: daily report | Example: live dashboard | +-----------------------+-----------------------+
📌 Mini summary: Batch processes data all at once; stream processes data as it arrives.
Definition: Spark Streaming is a module in Spark that allows us to process streaming data in real-time.
Why it is important: It brings the power of Spark to live data, allowing us to analyse tweets, sensor data, transactions, and more in real-time.
Simple explanation: Spark Streaming is like a live filter that catches and processes data as it flies by.
Real-life example: Uber uses Spark Streaming to track ride requests and driver locations in real-time.
School example: A system that counts students as they enter the school compound.
Home example: A smart doorbell that sends data every time someone rings it.
Nigerian example: A traffic monitoring system that counts cars on the Third Mainland Bridge in Lagos.
Spark Streaming Architecture:
+-------------+ +-------------+ +-------------+
| Data Source | → | Spark | → | Output |
| (Kafka, | | Streaming | | (Dashboard, |
| Socket, etc)| | Engine | | Database, |
+-------------+ +-------------+ | Alert) |
+-------------+
📌 Mini summary: Spark Streaming is Spark's tool for processing live data as it arrives.
Definition: Spark Streaming processes streaming data by dividing it into small chunks called micro‑batches. Each micro‑batch is processed like a small batch job.
Why it is important: This approach combines the best of batch and streaming. It's fast and reliable.
Simple explanation: Imagine you are drinking juice through a straw. You don't drink it all at once. You take small sips. Each sip is a micro‑batch.
Real-life example: A bus picks up passengers every 10 minutes. Each busload is a micro‑batch.
School example: The teacher collects homework every morning. Each morning's collection is a micro‑batch.
Home example: You check your phone every few minutes for new messages. Each check is a micro‑batch.
Nigerian example: A market trader counts sales every hour. Each hour is a micro‑batch.
Streaming Data Flow: Time: |----|----|----|----|----|----| Data: | d1 | d2 | d3 | d4 | d5 | d6 | Micro-batches: [d1,d2] [d3,d4] [d5,d6] Spark processes each micro-batch one at a time.
📌 Mini summary: Spark Streaming uses micro‑batches – small chunks of data processed one at a time.
Definition: DStream stands for Discretized Stream. It is the basic data structure in Spark Streaming, representing a continuous stream of data.
Why it is important: DStreams are to streaming what RDDs are to batch – the fundamental building block.
Simple explanation: A DStream is like a conveyor belt with boxes of data. Each box is a small RDD.
Real-life example: A conveyor belt at an airport luggage carousel. Each bag is a piece of data.
School example: A line of students entering the assembly hall. Each student is a piece of data.
Home example: A stream of water droplets from a dripping tap. Each droplet is data.
Nigerian example: A line of people waiting at a bus stop. Each person is a piece of data.
DStream = continuous sequence of RDDs Time: |----|----|----|----| RDDs: [RDD1][RDD2][RDD3][RDD4] Each RDD contains data from one micro-batch interval.
📌 Mini summary: DStream is a stream of data divided into small RDDs.
Definition: A StreamingContext is the main entry point for Spark Streaming applications, similar to SparkContext for batch.
Why it is important: You need a StreamingContext to create and manage your streaming pipeline.
Simple explanation: It's like the control centre for your streaming application.
Real-life example: The air traffic control tower that manages all flights.
School example: The principal's office that manages school activities.
Home example: The remote control for your TV – it controls everything.
Nigerian example: The dispatcher at a taxi company that coordinates all drivers.
How to create a StreamingContext:
from pyspark import SparkContext
from pyspark.streaming import StreamingContext
sc = SparkContext("local[2]", "StreamingApp")
ssc = StreamingContext(sc, batchDuration=1) # 1 second micro-batches
# Now you can create DStreams and process data
📌 Mini summary: StreamingContext is the main object for creating Spark Streaming applications.
Definition: Data sources are places where streaming data originates, like message queues, sockets, or files.
Why it is important: You need to connect to a data source to receive streaming data.
Simple explanation: It's like a pipe bringing water into your house. You need to connect to it.
Real-life example: Kafka is a popular data source that streams messages between applications.
School example: The school bell is a source – it sends a signal every hour.
Home example: A smart doorbell sends data when someone rings it.
Nigerian example: A mobile app sends data when a user makes a transaction.
Common Streaming Sources:
1. Socket (network connection)
2. Kafka (message queue)
3. Files (new files in a directory)
4. Flume (log collector)
5. Kinesis (AWS streaming service)
Example: reading from a socket
lines = ssc.socketTextStream("localhost", 9999)
📌 Mini summary: Data sources like Kafka or sockets provide the streaming data to Spark.
Definition: Transformations are operations that change a DStream, similar to transformations on RDDs.
Why it is important: Transformations let you clean, filter, and prepare streaming data.
Simple explanation: It's like having a sieve that separates sand from stones as the water flows through.
Real-life example: Filtering out all tweets that are not in English.
School example: Separating students by grade as they enter the school.
Home example: Sorting laundry into whites and colours as you pick it up.
Nigerian example: Filtering transactions by amount (> ₦10,000).
Common DStream Transformations:
- map() : apply a function to each element
- filter() : keep only elements that satisfy a condition
- flatMap() : split each element into multiple elements
- union() : combine two DStreams
- reduceByKey(): aggregate by key
- count() : count elements
Example: Word count on streaming data
words = lines.flatMap(lambda line: line.split(" "))
pairs = words.map(lambda word: (word, 1))
word_counts = pairs.reduceByKey(lambda a, b: a + b)
📌 Mini summary: Transformations on DStreams change, filter, and process the streaming data.
Definition: Actions are operations that produce results from a DStream, like printing or saving data.
Why it is important: Actions are what you do with the processed data – you need to see results!
Simple explanation: After you've sorted the laundry, you put it in the drawers – that's the action.
Real-life example: A dashboard that displays live sales data.
School example: A bell that rings when a certain number of students arrive.
Home example: A notification on your phone when you get a message.
Nigerian example: An SMS alert when a large transaction is made.
Common DStream Actions:
- print() : print the first 10 elements
- saveAsTextFiles(): save to text files
- saveAsHadoopFiles(): save to Hadoop
- count() : count elements (returns a stream)
- foreachRDD() : apply a function to each RDD (most flexible)
Example: Print word counts
word_counts.print()
Example: Save to file
word_counts.saveAsTextFiles("output/wordcount")
📌 Mini summary: Actions output the results of your streaming processing.
Definition: Window operations allow you to process data over a sliding window of time, not just the current micro‑batch.
Why it is important: Sometimes you want to see trends over the last 5 minutes, not just the last second.
Simple explanation: It's like looking through a moving window. You see what's happening now and what happened recently.
Real-life example: A traffic report shows traffic conditions over the last 15 minutes.
School example: A teacher looks at the last 5 minutes of students entering to estimate total attendance.
Home example: You check how many messages you received in the last hour.
Nigerian example: A security monitor checks the last 10 minutes of camera footage.
Window Operations:
+------------------------------------------+
| windowDuration = how long the window is |
| slideDuration = how often the window moves|
+------------------------------------------+
Example: Count words in the last 10 seconds, updated every 2 seconds
windowed_counts = pairs.reduceByKeyAndWindow(
lambda a, b: a + b, # add function
lambda a, b: a - b, # subtract function (for sliding)
10, # window duration (10 seconds)
2 # slide duration (2 seconds)
)
📌 Mini summary: Window operations let you analyse data over a period of time, like the last 10 seconds.
Definition: Stateful operations maintain information across micro‑batches. They "remember" what happened before.
Why it is important: Some analyses need to accumulate data over time, like total visitors per day.
Simple explanation: It's like keeping a running total of how many goals have been scored in a match so far.
Real-life example: Tracking total sales for the day, updated every minute.
School example: Keeping a running count of total students who have arrived in the morning.
Home example: Keeping a running total of how many steps you've walked today.
Nigerian example: Keeping a running count of total votes cast in an election as they arrive.
Stateful Operations:
1. updateStateByKey() - maintains state for each key
2. mapWithState() - more efficient state management
Example: Running count per word
def update_func(new_values, running_count):
if running_count is None:
running_count = 0
return sum(new_values) + running_count
running_counts = pairs.updateStateByKey(update_func)
📌 Mini summary: Stateful operations remember data across micro‑batches.
Definition: Late data is data that arrives after its expected time. Watermarks help the system decide when to wait for late data.
Why it is important: In real-world streaming, data can arrive late due to network delays or other issues.
Simple explanation: It's like waiting a little bit longer for a friend who is running late. You set a time you'll wait (the watermark).
Real-life example: A sensor sends data but there is a network delay.
School example: A student arrives late to class – you still count them but you know they are late.
Home example: A package arrives a day late – you still accept it.
Nigerian example: Election results from a remote village arrive late – they are still counted.
Watermark in Spark Structured Streaming:
from pyspark.sql.functions import current_timestamp
streaming_df = spark.readStream.format("kafka")...
# Add watermark (wait 10 minutes for late data)
streaming_df = streaming_df.withWatermark("timestamp", "10 minutes")
# Now window operations will handle late data
📌 Mini summary: Watermarks handle late data by setting how long to wait for it.
Definition: Structured Streaming is the newer, more powerful version of Spark Streaming. It uses DataFrames and SQL instead of DStreams.
Why it is important: It's easier to use, more powerful, and handles exactly‑once processing.
Simple explanation: It's like upgrading from a bicycle to a car – faster, smoother, and more features.
Real-life example: Netflix uses Structured Streaming to analyse viewing data in real-time.
School example: A school uses a modern attendance system that updates live.
Home example: A smart home system that monitors all devices in real-time.
Nigerian example: A modern payment system that processes transactions in real-time.
Structured Streaming Example:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("StructuredStreaming").getOrCreate()
# Read streaming data
lines = spark.readStream.format("socket")
.option("host", "localhost")
.option("port", 9999)
.load()
# Process with SQL
words = lines.selectExpr("split(value, ' ') as words")
word_counts = words.groupBy("words").count()
# Write the output (sink)
query = word_counts.writeStream.outputMode("complete")
.format("console")
.start()
query.awaitTermination()
📌 Mini summary: Structured Streaming is the newer, easier way to do streaming with DataFrames.
Definition: Output sinks are where the results of streaming processing are sent, like a database, file, or console.
Why it is important: You need to deliver your results somewhere useful.
Simple explanation: It's like deciding where to put the food after cooking – on a plate, in a box, or in the fridge.
Real-life example: A dashboard displays results on a screen.
School example: Results are written on the notice board.
Home example: A notification pops up on your phone.
Nigerian example: Results are saved to a database for later analysis.
Output Sinks in Structured Streaming:
1. Console sink: print to the console (good for testing)
2. File sink: save to files (CSV, Parquet, etc.)
3. Kafka sink: send to Kafka topic
4. Foreach sink: custom processing (write to database, etc.)
5. Memory sink: store in memory (for testing)
Example: Write to console
query = word_counts.writeStream.outputMode("complete")
.format("console")
.start()
Example: Write to Parquet files
query = word_counts.writeStream.outputMode("append")
.format("parquet")
.option("path", "output/")
.start()
📌 Mini summary: Output sinks decide where to send the processed streaming results.
Let's build a complete real‑time dashboard that shows live sales data from a bakery.
Problem: A bakery wants to see live sales: total items sold, most popular item, and revenue.
Data Source: A socket sends sales data: item, price, quantity, timestamp.
Steps:
# Import Spark
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
spark = SparkSession.builder.appName("BakeryDashboard").getOrCreate()
# 1. Read streaming data from socket
sales_df = spark.readStream.format("socket")
.option("host", "localhost")
.option("port", 9999)
.load()
# 2. Parse the data (assuming format: item,price,quantity)
sales_df = sales_df.selectExpr(
"split(value, ',')[0] as item",
"cast(split(value, ',')[1] as double) as price",
"cast(split(value, ',')[2] as int) as quantity"
)
# 3. Add total_sales column
sales_df = sales_df.withColumn("total_sales", sales_df.price * sales_df.quantity)
# 4. Window operations - 1 minute windows
sales_df = sales_df.withWatermark("timestamp", "1 minute")
# 5. Calculate aggregations
totals = sales_df.groupBy("item").agg(
sum("total_sales").alias("revenue"),
sum("quantity").alias("items_sold")
)
# 6. Write to console (dashboard)
query = totals.writeStream.outputMode("complete")
.format("console")
.trigger(processingTime="5 seconds")
.start()
# 7. Also write to memory for dashboard
query2 = totals.writeStream.outputMode("complete")
.format("memory")
.queryName("dashboard")
.start()
query.awaitTermination()
📌 Mini summary: A real‑time dashboard shows live data updates as they happen.
How to build a Spark Streaming application:
Step 1: Create StreamingContext Step 2: Connect to source (socket, Kafka) Step 3: Process with transformations Step 4: Apply windows/state Step 5: Write to sink Step 6: Start streaming Step 7: Wait/stop
Streaming Pipeline:
Data Source → Spark Streaming → Transformations → Window/State → Sink → Results (Kafka) (micro-batches) (map, filter) (time windows) (DB) (Dashboard)
Micro‑batch Concept:
Data Stream: |------|------|------|------|------|------| | d1d2 | d3d4 | d5d6 | d7d8 | d9d10| ... | |----| |----| |----| |----| |----| |----| Batch1 Batch2 Batch3 Batch4 Batch5 Each batch = 2 seconds of data
Window Operation:
Time: |--|--|--|--|--|--|--|--|--|
Windows: |---------| (10 seconds)
|---------| (slides every 2 seconds)
|---------|
Each window covers the last 10 seconds
Structured Streaming Flow:
Input Stream → DataFrame → SQL/Transform → Output Stream (live data) (table) (groupBy, etc) (console/file/DB)
Batch vs Stream Processing:
| Feature | Batch | Stream |
|---|---|---|
| Data | Static/Complete | Continuous/Infinite |
| Processing | All at once | As it arrives |
| Latency | Hours to days | Milliseconds to seconds |
| Use case | Daily reports | Live dashboards |
| Output | Complete result | Continuous updates |
DStream vs Structured Streaming:
| Feature | DStream (Old) | Structured Streaming (New) |
|---|---|---|
| API | RDD-based | DataFrame/SQL-based |
| Ease of use | Harder | Easier |
| Exactly-once | Harder | Built-in |
| Late data | Manual | Watermarks |
| Performance | Good | Better |
Streaming Sources Comparison:
| Source | Description | Best For |
|---|---|---|
| Socket | Network connection | Testing |
| Kafka | Message queue | Production |
| Files | New files in a directory | Log files |
| Kinesis | AWS streaming | AWS users |
Congratulations! 🎉 You have completed Module 7 – Spark Streaming and Real‑Time Data. Here's what we learned:
You are now ready for Module 8 – Machine Learning with Spark MLlib! 🤖
Match the term with its description:
| Term | Description |
|---|---|
| 1. Stream | A. Small chunk of streaming data |
| 2. Micro‑batch | B. Data that arrives continuously |
| 3. DStream | C. Processing over time |
| 4. Window | D. Handles late data |
| 5. Watermark | E. Streaming data structure |
Answers: 1-B, 2-A, 3-E, 4-C, 5-D
Title: Design a Real‑Time Monitoring System
Instructions: In groups of 4, design a real‑time monitoring system for a Nigerian business of your choice (e.g., a bank, a supermarket, a traffic system). Answer:
Title: Simple Socket Streaming
Instructions: Set up a simple socket server on your computer. Write a Spark Streaming application that reads from the socket and counts words. Send some sentences to the socket and observe the word counts updating live.
Title: Real‑Time Social Media Monitoring
Description: Build a Spark Streaming application that monitors social media posts (simulated by a socket). The application should:
Deliverable: Working Spark Streaming code.
Title: Kafka and Spark Streaming
Instructions:
Title: Real‑Time Anomaly Detection
Problem: A sensor network sends data every second. Build a streaming application that detects anomalies (values that are more than 3 standard deviations from the mean). Use a sliding window to compute the mean and standard deviation in real‑time.
Fill‑in‑the‑Blank Answers: 1. continuously, 2. all at once, 3. micro‑batches, 4. Discretized Stream, 5. DataFrames
True/False Answers: 1. True, 2. False, 3. True, 4. False, 5. True
Multiple Choice Answers: 1-B, 2-B, 3-B, 4-B, 5-B, 6-B, 7-B, 8-B, 9-B, 10-B, 11-B, 12-C, 13-B, 14-B, 15-A
Matching Answers: 1-B, 2-A, 3-E, 4-C, 5-D
In Module 8, we will explore the exciting world of Machine Learning with Spark. We will learn how to build models that can learn from data and make predictions.
What to bring:
See you in Module 8 – let's teach computers to learn! 🤖
Module Introduction
Hello, future AI engineer! 🧠 In Module 7, we learned how to process data in real-time with Spark Streaming. Now, we are going to teach computers how to learn from data – this is called Machine Learning!
Imagine you have a magic computer that can learn from examples. You show it pictures of cats and dogs. It looks at them and learns what makes a cat a cat and a dog a dog. Then, when you show it a new picture, it can tell you if it's a cat or a dog. That's machine learning!
Spark has a special library called MLlib (Machine Learning Library) that lets us do machine learning on big data. MLlib is like a treasure chest full of algorithms (recipes) that computers can use to learn from data.
In this module, we will learn about the basics of machine learning, different types of learning, and how to use MLlib to build models that can predict, classify, and find patterns. Let's begin! 🚀
Learning Objectives
By the end of this module, you will be able to:
Warm‑up Story – The Smart Farmer
Once upon a time, in a village in Nigeria, there was a farmer named Chief Ade. Chief Ade had a big farm with many yams. Every year, he had to decide how many yams to plant, how much water to give them, and when to harvest them.
Chief Ade had old notebooks with data from the past 50 years. He noticed patterns: when it rained a lot, yams grew bigger. When it was too hot, they grew smaller. But there were too many factors to keep track of in his head!
His grandson, Femi, was learning about computers. He said, "Grandpa, let's use your data to teach a computer to predict your yam harvest! The computer will look at all your old data and learn the patterns. Then, when you tell it how much rain is expected, it will predict how many yams you will get!"
Chief Ade was amazed. They used a machine learning model. It looked at rainfall, temperature, and soil quality from the past. The model learned the relationship between these factors and the yam harvest. Now, Chief Ade could predict his harvest before planting! He could plan better and make more money.
That is exactly what machine learning does – it finds patterns in data and makes predictions.
Old Data → Machine Learning → Predictions (rain, temp, soil) (model) (yam harvest)
🌟 Mini summary: Machine learning teaches computers to find patterns in data and make predictions.
Definition: Machine learning is a way for computers to learn from data without being explicitly programmed. Instead of writing rules, we give the computer data and let it discover patterns on its own.
Why it is important: Machine learning helps us make predictions and decisions from large amounts of data.
Simple explanation: It's like teaching a child by showing them many examples. The child learns the pattern without you telling them every rule.
Real-life example: Your phone's keyboard learns which words you use and predicts what you will type next.
School example: You learn to recognise different animals by looking at many pictures. That's machine learning!
Home example: A smart thermostat learns your schedule and adjusts the temperature accordingly.
Nigerian example: A bank uses machine learning to decide if a customer is likely to pay back a loan.
Traditional Programming: Rules + Data → Answers (you tell the computer what to do) Machine Learning: Data + Answers → Rules (the computer learns the rules)
📌 Mini summary: Machine learning lets computers learn patterns from data and make predictions.
Definition: Supervised learning is when we train a model using labelled data – data that has the correct answer already known.
Why it is important: This is the most common type of machine learning. It's used for predictions.
Simple explanation: It's like having a teacher who gives you the correct answers so you can learn.
Real-life example: A spam filter learns from emails that are labelled "spam" or "not spam".
School example: Your teacher gives you practice tests with answers. You learn from them.
Home example: Your parents show you examples of good behaviour and bad behaviour.
Nigerian example: A hospital uses labelled patient data to predict if someone has malaria.
Supervised Learning: Input (X) + Labels (Y) → Model → Predictions (features) (answers) (learn) (for new data) Types: - Classification: Predict a category (cat or dog) - Regression: Predict a number (price of yams)
📌 Mini summary: Supervised learning uses data with correct answers to train a model.
Definition: Unsupervised learning is when we have data but no correct answers. The model tries to find patterns and groups in the data on its own.
Why it is important: Sometimes we don't have labelled data. Unsupervised learning helps us discover hidden patterns.
Simple explanation: It's like sorting a box of mixed toys without knowing the categories. You decide how to group them based on similarities.
Real-life example: An online store groups customers into segments based on their shopping behaviour.
School example: You group your classmates by their favourite subjects without being told.
Home example: You organise your toys by colour, size, or type without anyone telling you how.
Nigerian example: A bank groups customers by their spending habits to offer them special deals.
Unsupervised Learning: Input (X) → Model → Patterns/Groups (features) (learn) (clusters) Types: - Clustering: Group similar data together - Dimensionality Reduction: Simplify data
📌 Mini summary: Unsupervised learning finds patterns in data without any correct answers.
Definition: MLlib is Spark's machine learning library. It provides tools and algorithms for building machine learning models on big data.
Why it is important: MLlib makes it easy to do machine learning on large datasets. It's fast and scales across many computers.
Simple explanation: MLlib is like a big box of building blocks. You can use these blocks to build your own machine learning applications.
Real-life example: A company uses MLlib to build a recommendation system for its e‑commerce website.
School example: Your teacher uses a tool to predict which students might need extra help.
Home example: A smart home system learns your preferences and adjusts settings.
Nigerian example: A Nigerian startup uses MLlib to predict crop yields from weather data.
MLlib Components: +------------------+---------------------------+ | Component | What it does | +------------------+---------------------------+ | ML Algorithms | Classification, regression| | Feature Tools | Prepare data for learning | | Pipelines | Chain steps together | | Evaluation | Measure model performance | | Persistence | Save and load models | +------------------+---------------------------+
📌 Mini summary: MLlib is Spark's machine learning library that helps us build ML models on big data.
Definition: A pipeline is a sequence of steps that take raw data and produce a trained model. It includes data preparation, feature engineering, and model training.
Why it is important: Pipelines organise the whole machine learning process. They make it repeatable and clean.
Simple explanation: It's like a factory assembly line – raw materials (data) go in one end and finished products (predictions) come out the other.
Real-life example: A car factory has an assembly line with many steps. Each step does one part of the job.
School example: You follow a recipe step by step to bake a cake.
Home example: A dishwasher has a cycle: rinse, wash, dry.
Nigerian example: A factory processes cocoa beans: cleaning, roasting, grinding, packaging.
ML Pipeline: Raw Data → Clean Data → Feature Engineering → Model Training → Evaluation → Model (messy) (remove bad) (create features) (learn from data) (test) (ready) In MLlib: from pyspark.ml import Pipeline pipeline = Pipeline(stages=[stage1, stage2, stage3, ...]) model = pipeline.fit(train_data)
📌 Mini summary: A pipeline organises all the steps from raw data to a trained model.
Definition: Feature engineering is the process of turning raw data into features (inputs) that a machine learning model can understand.
Why it is important: Good features = good predictions. Feature engineering is one of the most important parts of machine learning.
Simple explanation: It's like preparing ingredients before cooking. You wash, peel, and chop the vegetables before using them in the recipe.
Real-life example: In a house price prediction model, features could be: number of rooms, size, location, age of the house.
School example: Features for predicting student performance: hours studied, sleep hours, attendance.
Home example: Features for predicting energy usage: number of people, house size, outside temperature.
Nigerian example: Features for predicting election results: previous votes, location, voter turnout.
Feature Engineering Steps: 1. Handle missing values (fill with average or delete) 2. Convert text to numbers (using StringIndexer) 3. Normalize numbers (scale to similar ranges) 4. Create new features from existing ones 5. Select the most important features MLlib Example: from pyspark.ml.feature import VectorAssembler assembler = VectorAssembler(inputCols=["feature1", "feature2"], outputCol="features") data_with_features = assembler.transform(data)
📌 Mini summary: Feature engineering prepares data so that machine learning models can use it.
Definition: Classification is a type of supervised learning where we predict a category or class.
Why it is important: Many problems involve categories: spam or not spam, sick or healthy, good or bad.
Simple explanation: It's like sorting objects into boxes. Each box has a label.
Real-life example: A bank classifies credit card transactions as "fraudulent" or "legitimate".
School example: A teacher classifies test scores as "pass" or "fail".
Home example: You classify clothes as "clean" or "dirty".
Nigerian example: A hospital classifies malaria test results as "positive" or "negative".
Classification Algorithms in MLlib: 1. Logistic Regression 2. Decision Trees 3. Random Forest 4. Naive Bayes 5. Support Vector Machines (SVM) Example: Predict if a customer will buy a product (Yes/No) from pyspark.ml.classification import RandomForestClassifier rf = RandomForestClassifier(labelCol="label", featuresCol="features") model = rf.fit(train_data)
📌 Mini summary: Classification predicts which category something belongs to.
Definition: Regression is a type of supervised learning where we predict a continuous number (not a category).
Why it is important: Many problems involve numbers: price, temperature, sales, age.
Simple explanation: It's like guessing a number within a range.
Real-life example: Predicting the price of a house based on its features.
School example: Predicting a student's final exam score based on their study hours.
Home example: Predicting how much your electricity bill will be based on usage.
Nigerian example: Predicting yam prices based on rainfall and temperature.
Regression Algorithms in MLlib: 1. Linear Regression 2. Decision Tree Regression 3. Random Forest Regression 4. Gradient Boosted Trees Example: Predict house price from pyspark.ml.regression import LinearRegression lr = LinearRegression(labelCol="label", featuresCol="features") model = lr.fit(train_data)
📌 Mini summary: Regression predicts a continuous number value.
Definition: Clustering is an unsupervised learning technique that groups similar data points together.
Why it is important: Clustering helps us discover natural groups in data that we didn't know existed.
Simple explanation: It's like sorting a box of mixed candies by colour and size, without any labels.
Real-life example: An online store groups customers by shopping behaviour (frequent buyers, occasional buyers, etc.).
School example: Grouping students by their favourite subjects.
Home example: Grouping your toys by type (cars, dolls, puzzles).
Nigerian example: Grouping farmers by the type of crops they grow.
Clustering Algorithms in MLlib: 1. K-Means 2. Bisecting K-Means 3. Gaussian Mixture Model (GMM) Example: Group customers by spending habits from pyspark.ml.clustering import KMeans kmeans = KMeans(k=3, featuresCol="features") model = kmeans.fit(data) predictions = model.transform(data)
📌 Mini summary: Clustering finds natural groups in data without any labels.
Definition: ALS (Alternating Least Squares) is a collaborative filtering algorithm used for recommendation systems.
Why it is important: Recommendation systems help users discover new products, movies, or content they might like.
Simple explanation: It finds patterns in what users like and recommends new things based on those patterns.
Real-life example: Netflix recommends movies you might like based on what you've watched.
School example: A teacher recommends books based on what students have read before.
Home example: YouTube recommends videos based on your viewing history.
Nigerian example: Jumia recommends products based on what you've bought before.
ALS Recommendation System:
User Ratings → ALS Model → Predictions → Recommendations
(user, item, rating) (learn) (predict) (top items)
Example: Movie recommendations
from pyspark.ml.recommendation import ALS
als = ALS(maxIter=5, regParam=0.01, userCol="userId",
itemCol="movieId", ratingCol="rating",
nonnegative=True)
model = als.fit(ratings_df)
recommendations = model.recommendForAllUsers(5)
📌 Mini summary: ALS is a recommendation algorithm that suggests items based on user preferences.
Definition: Evaluation means measuring how well your machine learning model performs.
Why it is important: You need to know if your model is good enough to use. Evaluation tells you if it's working.
Simple explanation: It's like checking your homework to see if you got the answers right.
Real-life example: A bank tests its fraud detection model to see how many frauds it catches.
School example: Your teacher marks your test to evaluate your performance.
Home example: You taste your food to see if it's good.
Nigerian example: A hospital tests a malaria detection model to see how accurate it is.
Evaluation Metrics:
For Classification:
- Accuracy: % of correct predictions
- Precision: % of positive predictions that were correct
- Recall: % of actual positives that were found
- F1 Score: balance between precision and recall
For Regression:
- RMSE (Root Mean Squared Error)
- R² (R-squared)
MLlib Example:
from pyspark.ml.evaluation import MulticlassClassificationEvaluator
evaluator = MulticlassClassificationEvaluator(labelCol="label",
predictionCol="prediction",
metricName="accuracy")
accuracy = evaluator.evaluate(predictions)
📌 Mini summary: Evaluation tells us how well our machine learning model is performing.
Definition: Train-test split means dividing your data into two parts: one for training the model and one for testing it.
Why it is important: If we test on the same data we trained on, we might be fooled. The model might have memorised the answers instead of learning patterns. This is called overfitting.
Simple explanation: It's like studying for a test. If you study the exact questions that will be on the test, you'll do well – but you didn't really learn. A proper test should have new questions!
Real-life example: A teacher uses practice tests (training) and final exams (testing) separately.
School example: You do homework (training) and then take a test with new questions (testing).
Home example: You practice cooking one dish and then cook it for guests (testing).
Nigerian example: A model predicts election results using past data (training) and is tested on future data (testing).
Train-Test Split: +------------------+------------------+ | Training Data | Testing Data | | (70-80% of data) | (20-30% of data) | | Used to train | Used to evaluate | | the model | the model | +------------------+------------------+ MLlib Example: train_data, test_data = data.randomSplit([0.8, 0.2]) model = algorithm.fit(train_data) predictions = model.transform(test_data) accuracy = evaluator.evaluate(predictions)
📌 Mini summary: Train-test split helps us evaluate our model properly by testing it on unseen data.
Definition: Feature transformers are MLlib tools that convert data from one format to another for machine learning.
Why it is important: Machine learning algorithms need data in specific formats (numerical vectors). Transformers convert data into those formats.
Simple explanation: It's like a translator that converts your data into a language the ML model can understand.
Real-life example: Converting text ("red", "blue") into numbers (1, 2) using StringIndexer.
School example: Changing grades (A, B, C) into numbers (4, 3, 2).
Home example: Converting temperature from Celsius to Fahrenheit.
Nigerian example: Converting Nigerian states (Lagos, Abuja, Kano) into numbers for a model.
Common MLlib Transformers:
1. StringIndexer: Convert categories to numbers
2. VectorAssembler: Combine columns into a vector
3. StandardScaler: Normalize numbers (scale to similar range)
4. OneHotEncoder: Convert categories to binary vectors
5. Tokenizer: Split text into words
Example:
from pyspark.ml.feature import StringIndexer, VectorAssembler
indexer = StringIndexer(inputCol="city", outputCol="cityIndex")
indexed = indexer.fit(data).transform(data)
assembler = VectorAssembler(inputCols=["cityIndex", "age"],
outputCol="features")
ready_data = assembler.transform(indexed)
📌 Mini summary: Feature transformers convert data into formats that machine learning algorithms can use.
Definition: Saving and loading models means storing a trained model on disk so we can use it later without retraining.
Why it is important: Training models can take a long time. We don't want to retrain every time we want to make a prediction.
Simple explanation: It's like saving your game progress so you don't have to start from the beginning every time.
Real-life example: A model trained on months of data is saved and used to make predictions every day.
School example: You save your project so you don't lose your work.
Home example: You save your favourite recipe so you can use it again.
Nigerian example: A bank saves a fraud detection model and uses it for every transaction.
Saving and Loading Models:
# Save a model
model.save("path/to/model")
# Load a model
from pyspark.ml.classification import RandomForestClassificationModel
loaded_model = RandomForestClassificationModel.load("path/to/model")
# Use loaded model for predictions
predictions = loaded_model.transform(new_data)
📌 Mini summary: Saving and loading models lets us reuse trained models without retraining.
Let's build a complete machine learning project: Predicting Student Performance
Problem: A school wants to predict which students might need extra help based on their data.
Data: Student records: hours_studied, attendance, previous_score, sleep_hours, and final_score (label).
Complete Code:
from pyspark.sql import SparkSession
from pyspark.ml.feature import VectorAssembler, StandardScaler
from pyspark.ml.regression import LinearRegression
from pyspark.ml.evaluation import RegressionEvaluator
spark = SparkSession.builder.appName("StudentPerformance").getOrCreate()
# Step 1: Load data
data = spark.read.csv("student_performance.csv", header=True, inferSchema=True)
# Step 2: Prepare features
feature_columns = ["hours_studied", "attendance", "previous_score", "sleep_hours"]
assembler = VectorAssembler(inputCols=feature_columns, outputCol="features_raw")
data = assembler.transform(data)
# Step 3: Scale features
scaler = StandardScaler(inputCol="features_raw", outputCol="features",
withStd=True, withMean=True)
scaler_model = scaler.fit(data)
data = scaler_model.transform(data)
# Step 4: Select final data
final_data = data.select("features", "final_score")
final_data = final_data.withColumnRenamed("final_score", "label")
# Step 5: Train-test split
train_data, test_data = final_data.randomSplit([0.8, 0.2], seed=42)
# Step 6: Train model
lr = LinearRegression(labelCol="label", featuresCol="features")
model = lr.fit(train_data)
# Step 7: Evaluate model
predictions = model.transform(test_data)
evaluator = RegressionEvaluator(labelCol="label", predictionCol="prediction",
metricName="rmse")
rmse = evaluator.evaluate(predictions)
print(f"RMSE on test data: {rmse}")
# Step 8: Show predictions
predictions.select("label", "prediction").show(10)
# Step 9: Save model
model.save("student_performance_model")
print("Model saved! 🎉")
Sample Output: +-------+------------------+ | label | prediction | +-------+------------------+ | 85.0 | 84.7 | | 70.0 | 71.2 | | 92.0 | 91.5 | | 63.0 | 62.8 | +-------+------------------+ RMSE: 2.34
📌 Mini summary: A complete ML project takes data through preparation, training, evaluation, and saving.
How to build an ML model with Spark MLlib:
Step 1: Load Data Step 2: Prepare Data Step 3: Define Features & Label Step 4: Apply Feature Transformers Step 5: Train-Test Split Step 6: Create Algorithm Step 7: Train Model Step 8: Make Predictions Step 9: Evaluate Model Step 10: Save Model
ML Pipeline:
Raw Data → Transformer → Feature → Algorithm → Model → Evaluation → Predictions (messy) (clean/scale) (vector) (learn) (trained) (test) (results)
Supervised vs Unsupervised:
Supervised Learning (with teacher): +-------+-------+ +-------+ | Input | Label | → | Model | → Predictions +-------+-------+ +-------+ Unsupervised Learning (no teacher): +-------+ +-------+ | Input | → | Model | → Groups/Patterns +-------+ +-------+
Train-Test Split:
Total Data (100%) +------------------------------------------+ | Training Data (80%) | Testing Data (20%) | +------------------------------------------+ [Used to train model] [Used to evaluate]
Recommendation System Flow:
User Ratings → ALS Model → Predict Scores → Recommend Top Items (user, item, rating) (learn) (for each user) (top 5)
Supervised vs Unsupervised Learning:
| Feature | Supervised | Unsupervised |
|---|---|---|
| Labels | Yes (correct answers) | No |
| Goal | Predict answers | Find patterns |
| Examples | Classification, Regression | Clustering |
| Use case | Spam detection | Customer segmentation |
Classification vs Regression:
| Feature | Classification | Regression |
|---|---|---|
| Output | Category | Number |
| Examples | Spam/Not Spam | House Price |
| Algorithms | Random Forest | Linear Regression |
| Evaluation | Accuracy, F1 | RMSE, R² |
ML Algorithms in MLlib:
| Algorithm | Type | Use Case |
|---|---|---|
| Logistic Regression | Classification | Spam detection |
| Linear Regression | Regression | Price prediction |
| Random Forest | Both | Many problems |
| K-Means | Clustering | Customer segmentation |
| ALS | Recommendation | Movie recommendations |
Congratulations! 🎉 You have completed Module 8 – Machine Learning with Spark MLlib. Here's what we learned:
You now have the skills to build machine learning models on big data using Spark. Keep learning and building! 🤖
Match the term with its description:
| Term | Description |
|---|---|
| 1. Supervised | A. Predicts categories |
| 2. Unsupervised | B. Predicts numbers |
| 3. Classification | C. Learning with labels |
| 4. Regression | D. Learning without labels |
| 5. Clustering | E. Groups similar data |
Answers: 1-C, 2-D, 3-A, 4-B, 5-E
Title: Design a ML Product
Instructions: In groups of 4, design a machine learning product for a Nigerian problem. Answer:
Title: Build a Simple Classifier
Instructions: Use Spark MLlib to build a classification model. Create a small dataset with two features and a label. Train a Random Forest Classifier and evaluate its accuracy.
Title: Nigerian Food Recommendation System
Description: Build a recommendation system for Nigerian dishes. You have user ratings for different dishes. Use ALS to recommend dishes to users.
Deliverable: Spark code and recommendations for 5 users.
Title: Predict House Prices
Instructions:
Title: Improve Model Performance
Problem: Your model has an accuracy of 70%. Improve it to 85% using different algorithms, feature engineering, or hyperparameter tuning.
Fill‑in‑the‑Blank Answers: 1. learn, 2. labels, 3. categories, 4. numbers, 5. recommendations
True/False Answers: 1. False, 2. True, 3. False, 4. True, 5. True
Multiple Choice Answers: 1-B, 2-A, 3-B, 4-A, 5-B, 6-B, 7-A, 8-C, 9-B, 10-B, 11-B, 12-B, 13-B, 14-B, 15-A
Matching Answers: 1-C, 2-D, 3-A, 4-B, 5-E
Congratulations! You have completed all 8 modules of "Big Data Engineering with Spark"! 🎉
You have learned:
What's next in your big data journey:
Remember: The best way to learn is by doing. Keep building, keep experimenting, and never stop learning! You are now a data engineer – go and change the world with data! 🌍🚀