Big data Homeworks
Università degli Studi di Napoli Federico II
Autori: Giuseppe Laterza M63001411
Docente: Vincenzo Moscato
Anno Accademico: 2023/2024
Traccia I homework
(consegna 10/05/2024) - Dato uno o più dataset a libera scelta, implementare una serie di analytics relative a pattern di "filtering" e "summarization" con PySpark, utilizzando Google Colab e/o una piattaforma di cloud data management (es. Databricks, Kaggle, etc.).
Traccia II homework
(consegna 25/05/2024) - Dato uno o più dataset a libera scelta, implementare una serie di analytics relative a pattern di "filtering" e "summarization" con HIVE, utilizzando una piattaforma di cloud data management (es. Databricks, etc.).
Traccia III homework
(consegna giorno dell'esame) - Dato uno o più dataset a libera scelta, implementare l'import e memorizzazione dei dati ed una serie di analytics relative a pattern di "filtering" e "summarization" con un sistema nosql a scelta dello studente.
Introduzione
Lo scopo del primo Homework del corso di Big Data Engineering è quello di ricavare informazioni utili attraverso l’analisi di dati provenienti dall’ U.S. Department of Transportation’s Bureau of Transportation Statistics (BTS). Il lavoro svolto è in particolare incentrato sull’utilizzo di Apache Spark, un framework di elaborazione parallela open source che supporta l’elaborazione in memoria per migliorare le prestazioni delle applicazioni che analizzano Big Data.
Le soluzioni Big Data sono progettate per gestire insiemi troppo grandi o complessi di dati per i database tradizionali. Spark elabora grandi quantità di dati in memoria, molto più velocemente rispetto alle alternative.
Data characterization
Big data Homeworks 1
L’analisi si incentra sul dataset che contiene le informazioni dei voli interni agli Stati Uniti nell’anno 2015 e contiene oltre 4 milioni di record che descrivono in dettaglio le operazioni di volo, tra cui orari di partenza e arrivo, ritardi, cancellazioni e altro ancora.
Il dataset principale, denominato, è integrato con altri due file ( e ) per fornire flights.csv airlines.csv airports.csv un contesto più ampio sui nomi di compagnie aeree e aereoporti.
Campi chiave
- Informazioni temporali:, , , : Specificano l'anno, il mese, il giorno e il YEAR MONTH DAY DAY_OF_WEEK giorno della settimana del volo.
- Informazioni sul volo:: Codice della compagnia aerea. AIRLINE
- Aeroporti: : Codice dell'aeroporto di partenza. ORIGIN_AIRPORT
- : Codice dell'aeroporto di destinazione. DESTINATION_AIRPORT
- Orari e ritardi: , , : Orario programmato, SCHEDULED_DEPARTURE DEPARTURE_TIME DEPARTURE_DELAY orario effettivo e ritardo alla partenza (in minuti)., , : Orario programmato, orario SCHEDULED_ARRIVAL ARRIVAL_TIME ARRIVAL_DELAY effettivo e ritardo all'arrivo (in minuti).
- Cancellazioni e deviazioni:: Indica se il volo è stato cancellato (1 = sì, 0 = no). CANCELLED : Motivo della cancellazione (A = compagnia aerea, B CANCELLATION_REASON = condizioni meteorologiche, C = problemi di sicurezza, D = altro).
- : Indica se il volo è stato deviato (1 = sì, 0 = no). DIVERTED
- Informazioni sui ritardi specifici:, , , , AIR_SYSTEM_DELAY SECURITY_DELAY AIRLINE_DELAY LATE_AIRCRAFT_DELAY: Ritardi causati da diversi fattori (in minuti). WEATHER_DELAY
I dati all’interno della tabella sono stati sottoposti ad un preprocessing utilizzando PySpark. Queste operazioni sono state mirate a eliminare i valori Big data Homeworks 2 nulli rimpiazzandoli con altri.
Preprocessing dei dati
- Quando un volo è cancellato o deviato tutti i campi degli orari risultano nulli
- I campi relativi alle informazioni dei ritardi sono “null” se effettivamente non c’è ritardo, per questo motivo è stato inserito il campo ed DELEYED_FLIGHT impostati a 0 i valori nulli
- Il campo dei motivi della cancellazione era null se il volo non era cancellato. Sono stati sostituiti i valori nulli con la stringa “N” = Non cancellato
# Lista delle colonne da escludere dalla sostituzione
exclude_columns = ["AIR_SYSTEM_DELAY", "SECURITY_DELAY", "AIR
# Crea una lista delle colonne su cui applicare la trasformaz
columns_to_transform = [col for col in flights_df.columns if c
# Costruisci l'espressione per sostituire i null con 0 nelle c
for column in columns_to_transform:
print(column)
flights_df = flights_df.withColumn(column, when((((col("CA
# Sostituisci i valori null nella colonna CANCELLATION_REA
flights_df = flights_df.withColumn("CANCELLATION_REASON",when(col("CANCELLATION_REASON").isNull(), lit("N")).otherw)
# Aggiungi una colonna Flight_delayed che è 1 se almeno una def
flights_df = flights_df.withColumn("Flight_delayed",when((col("AIR_SYSTEM_DELAY").isNotNull()) |(col("SECURITY_DELAY").isNotNull()) |(col("AIRLINE_DELAY").isNotNull()) |(col("LATE_AIRCRAFT_DELAY").isNotNull()) |(col("WEATHER_DELAY").isNotNull()),1).otherwise(0)
Big data Homeworks 3
# Sostituisci i valori null con 0 nelle colonne di ritardo
delay_columns = ["AIR_SYSTEM_DELAY", "SECURITY_DELAY", "AIRLIN
for column in delay_columns:
flights_df = flights_df.withColumn(column,when(col(column).isNull(), 0).otherwise(col(column)))
Il risultato è quindi:
PySpark setup
In questa sezione viene descritto il processo di configurazione di Apache Spark in un ambiente Python, come ad esempio Google Colab. Innanzitutto Apache Spark è un framework open source per il calcolo distribuito su grandi quantità di dati, progettato per migliorare le prestazioni e la facilità d'uso rispetto a MapReduce. Spark è ideale per algoritmi di apprendimento automatico, analisi avanzate e interrogazioni complesse, sfruttando l'HDFS di Hadoop per l'archiviazione e la scalabilità su grandi volumi di dati.
Nel nostro caso è stato utilizzato PySpark, ovvero un interfaccia di Spark per Python.
Per funzionare Spark ha bisogno di Java, per questo motivo è stato installato Java 8. Una volta fatto questo viene scaricata ed estratta ed spark-3.4.2 installato findspark, una libreria Python che semplifica la configurazione di Spark in ambienti come Jupyter Notebook o Google Colab.
# Install Java 8
!apt-get install openjdk-8-jdk-headless -qq > /dev/null
Big data Homeworks 4
# Download Spark 3.4.2 (using the correct URL)
!wget https://archive.apache.org/dist/spark/spark-3.4.2/spark
# Extract Spark
!tar xf spark-3.4.2-bin-hadoop3.tgz
# Install findspark
!pip install findspark
Vengono poi configurate le variabili di ambiente necessarie a Spark per funzionare.
# Set up environment variables
import os
os.environ["JAVA_HOME"] = "/usr/lib/jvm/java-8-openjdk-amd64"
os.environ["SPARK_HOME"] = "/content/spark-3.4.2-bin-hadoop3"
# Initialize Spark
import findspark
findspark.init()
La SparkSession è il punto di ingresso principale per utilizzare le funzionalità di Spark, come l'API DataFrame e Spark SQL.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col
# Create Spark Session
spark = SparkSession.builder.appName("Flight Data Analysis").g
# Load the datasets
flights_df = spark.read.csv("/content/drive/MyDrive/flights.c
airlines_df = spark.read.csv("/content/drive/MyDrive/airlines
airports_df = spark.read.csv("/content/drive/MyDrive/airports
Big data Homeworks 5
In realtà Apache Spark è nato con gli RDD (Resilient Distributed Dataset), una struttura dati distribuita e immutabile, ma complessa da usare. Per semplificare il lavoro con dati strutturati, sono stati introdotti i DataFrame e successivamente i Dataset, che organizzano i dati in tabelle e permettono di scrivere query SQL-like. Proprio per la somiglianza con le tabelle è stato adottato questo approccio.
Analitica 1 - Statistiche generali sui voli
L’obiettivo di questa statistica è quello di avere delle informazioni generali sui dati come numero totale sui voli, range dei ritardi e tasso di cancellazione totale.
Inoltre è stato ritenuto interessante vedere come sono distribuite le partenze, poiché i ritardi all’arrivo derivano spesso da quest’ultime.
Per quanto riguarda la prima statistica sono state semplicemente calcolate facendo delle aggregazioni con funzioni di max e min. Si calcola il ritardo medio all'arrivo, avg_arr_delay, considerando solo i voli non cancellati e non deviati (filtrati con la condizione (CANCELLED == 0) & (DIVERTED == 0)).
# Statistiche di base su tutti i voli
def basic_statistics(df):
"""Calculate basic statistics about flights"""
return df.agg
count("*").alias("total_flights"),
min("DEPARTURE_DELAY").alias("min_dep_delay"),
max("DEPARTURE_DELAY").alias("max_dep_delay"),
avg("DEPARTURE_DELAY").alias("avg_dep_delay"),
min("ARRIVAL_DELAY").alias("min_arr_delay"),
max("ARRIVAL_DELAY").alias("max_arr_delay"),
# Calcola avg_arr_delay solo per voli non cancellati e
avg(when((col("CANCELLED") == 0) & (col("DIVERTED") ==
sum("CANCELLED").alias("total_cancellations"),
(sum("CANCELLED") / count("*") * 100).alias("cancella)
Big data Homeworks 6
Questa funzione analizza la distribuzione dei ritardi alla partenza, classificandoli in diverse categorie, come considerate dal BTS:
- very_early: Volo partito con un anticipo maggiore di 15 minuti.
- early : Volo partito con un anticipo compreso tra 1 e 15 minuti.
- on_time : Volo partito in orario o con un ritardo massimo di 15 minuti.
- slight_delay : Volo partito con un ritardo compreso tra 16 e 30 minuti.
- moderate_delay: Volo partito con un ritardo compreso tra 31 e 60 minuti.
- severe_delay: Volo partito con un ritardo superiore a 60 minuti.
Si utilizza la colonna DEPARTURE_DELAY per classificare i voli in base al ritardo.
Raggruppa i voli per categoria di ritardo ( delay_categ
Scarica il documento per vederlo tutto.
Scarica il documento per vederlo tutto.
Scarica il documento per vederlo tutto.
Scarica il documento per vederlo tutto.
Scarica il documento per vederlo tutto.
Scarica il documento per vederlo tutto.