Unit 1: Big data tools and analytics practice
Big Data Analytics Laboratory notes · PTU syllabus (PGCA1948)
On this page
Unit summary
This lab applies big data tools to real datasets: storing data in HDFS, running MapReduce and Spark jobs, querying with Hive, and implementing clustering, classification and association rule mining at scale.
After this unit you can
- Set up Hadoop and use HDFS commands
- Run MapReduce and Spark jobs
- Query data with Hive
- Implement clustering, classification and association rules with Spark MLlib
PTU syllabus topics
- Hands-on implementation of clustering
- classification and association-rule mining techniques using Hadoop/big data analytical tools on real-world datasets
- applying the concepts of HDFS
- MapReduce and the Hadoop ecosystem covered in the theory paper
HDFS
Distributed storage
YARN
Resource management
MapReduce
Batch processing
Hive
SQL-like queries
Pig
Data-flow scripts
HBase
NoSQL database
Spark
Fast in-memory processing
Topic 1
Setting up the environment
- 1Install Java 8/11 and Hadoop in pseudo-distributed mode (or use a Docker image or cloud cluster)
- 2Configure core-site.xml and hdfs-site.xml
- 3Format the NameNode (hdfs namenode -format)
- 4Start daemons (start-dfs.sh, start-yarn.sh) and check with jps
- 5Install PySpark (pip install pyspark)
- Web UIs: NameNode on port 9870, ResourceManager on port 8088.
Topic 2
HDFS commands
bashhdfs dfs -mkdir -p /lab/input
hdfs dfs -put retail.csv /lab/input/
hdfs dfs -ls /lab/input
hdfs dfs -cat /lab/input/retail.csv | head
hdfs dfs -get /lab/output/part-00000 result.txt
hdfs fsck /lab/input/retail.csv -files -blocks # see blocks and replicasTopic 3
MapReduce word count
bashhadoop jar $HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-examples-*.jar wordcount /lab/input /lab/wc_out
hdfs dfs -cat /lab/wc_out/part-r-00000 | sort -k2 -nr | head- Also run with Hadoop Streaming using Python mapper and reducer scripts.
Topic 4
Hive queries
sqlCREATE EXTERNAL TABLE retail (invoice STRING, item STRING, qty INT, price DOUBLE, country STRING)
ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' LOCATION '/lab/input' TBLPROPERTIES ("skip.header.line.count"="1");
SELECT country, ROUND(SUM(qty * price), 2) AS revenue FROM retail GROUP BY country ORDER BY revenue DESC LIMIT 10;Topic 5
Clustering with Spark MLlib
pythonfrom pyspark.sql import SparkSession
from pyspark.ml.feature import VectorAssembler, StandardScaler
from pyspark.ml.clustering import KMeans
from pyspark.ml.evaluation import ClusteringEvaluator
spark = SparkSession.builder.appName("lab").getOrCreate()
df = spark.read.csv("hdfs:///lab/input/customers.csv", header=True, inferSchema=True)
v = VectorAssembler(inputCols=["recency", "frequency", "monetary"], outputCol="raw").transform(df)
v = StandardScaler(inputCol="raw", outputCol="features").fit(v).transform(v)
for k in range(2, 7):
m = KMeans(k=k, seed=1).fit(v)
print(k, ClusteringEvaluator().evaluate(m.transform(v))) # silhouetteTopic 6
Classification with Spark MLlib
pythonfrom pyspark.ml.classification import DecisionTreeClassifier, NaiveBayes
from pyspark.ml.evaluation import MulticlassClassificationEvaluator
from pyspark.ml.feature import StringIndexer
data = spark.read.csv("hdfs:///lab/input/bank.csv", header=True, inferSchema=True)
data = StringIndexer(inputCol="y", outputCol="label").fit(data).transform(data)
data = VectorAssembler(inputCols=["age", "balance", "duration", "campaign"], outputCol="features").transform(data)
train, test = data.randomSplit([0.8, 0.2], seed=7)
ev = MulticlassClassificationEvaluator(metricName="accuracy")
for clf in (DecisionTreeClassifier(maxDepth=5), NaiveBayes()):
print(type(clf).__name__, ev.evaluate(clf.fit(train).transform(test)))Topic 7
Association rule mining with FP-Growth
pythonfrom pyspark.ml.fpm import FPGrowth
from pyspark.sql.functions import collect_set
baskets = spark.read.csv("hdfs:///lab/input/retail.csv", header=True).groupBy("invoice").agg(collect_set("item").alias("items"))
fp = FPGrowth(itemsCol="items", minSupport=0.02, minConfidence=0.3).fit(baskets)
fp.freqItemsets.orderBy("freq", ascending=False).show(10)
fp.associationRules.orderBy("lift", ascending=False).show(10) # antecedent, consequent, confidence, lift- FP-Growth avoids Apriori's repeated candidate generation and scans, so it scales better on large data.
Topic 8
Lab report checklist
- Dataset
- Source, size, fields
- Storage
- HDFS paths, block count
- Job
- Code, configuration, run time
- Results
- Tables and charts — clusters, accuracy, top rules
- Interpretation
- Business meaning and limitations
Key terms
- jps
- Command listing running Java (Hadoop) daemons
- External table
- Hive table over existing HDFS files
- VectorAssembler
- Combines columns into a feature vector in Spark
- FP-Growth
- Frequent-pattern mining without candidate generation
- Silhouette
- Cluster quality score
Quick revision
- Pseudo-distributed setup; daemons; web UIs.
- HDFS put, ls, cat, get, fsck.
- MapReduce word count; Hive external tables and aggregation.
- Spark k-means with silhouette; decision tree and Naive Bayes; FP-Growth with lift.
Important exam questions
Practice questions written to the PTU exam pattern for this unit's syllabus: short answers (Section A style) and long answers (Sections B and C style).
Short-answer questions
- Q1.How do you check that Hadoop daemons are running?
- Q2.Which command copies a file into HDFS?
- Q3.What is a Hive external table?
- Q4.What does VectorAssembler do?
- Q5.Why is FP-Growth preferred over Apriori for large data?
- Q6.Which metric evaluates k-means in Spark?
Long-answer questions
- Q1.Store a dataset in HDFS and run a MapReduce word count.
- Q2.Analyse sales data with Hive queries.
- Q3.Perform clustering and classification using Spark MLlib.
- Q4.Mine association rules from retail transactions with FP-Growth.
Stuck on this unit?
Message SBS on WhatsApp for help with Big Data Analytics Laboratory, or to ask about studying M.Sc IT at Synetic.
