# Define file location, file typem and CSV options file_location = "/mnt/adlstsdp/input/household_power_consumption.csv" file_type = "csv" schema = "Date STRING, Time STRING, Global_active_power DOUBLE, Global_reactive_power DOUBLE, Voltage DOUBLE, Global_intensity DOUBLE, Sub_metering_1 DOUBLE, Sub_metering_2 DOUBLE, Sub_metering_3 DOUBLE" first_row_is_header = "true" delimiter = ";" # Read CSV files org_df = spark.read.format(file_type) .schema(schema) .option("header", first_row_is_header) .option("delimiter", delimiter) .load(file_location) # Data cleansing and transformation from pyspark.sql.functions import * cleaned_df = org_df.na.drop() cleaned_df = cleaned_df.withColumn("Date", to_date(col("Date"),"d/M/y")) cleaned_df = cleaned_df.withColumn("Date", cleaned_df["Date"].cast("date")) cleaned_df = cleaned_df.select(concat_ws(" ", to_date(col("Date"),"d/M/y"), col("Time")).alias("DateTime"), "*") cleaned_df = cleaned_df.withColumn("DateTime", cleaned_df["DateTime"].cast("timestamp")) df = cleaned_df.groupby("Date").agg( round(sum("Global_active_power"), 2).alias("Total_global_active_power"), ).sort(["Date"]) # Add time-related features df = df.withColumn("year", year("Date")) df = df.withColumn("month", month("Date")) df = df.withColumn("week_num", weekofyear("Date")) # Add lagged value features of total global active power from pyspark.sql.window import Window from pyspark.sql.functions import lag windowSpec = Window.orderBy("Date") df = df.withColumn("power_lag1", round(lag(col("Total_global_active_power"), 1).over(windowSpec), 2)) # Create delta field df = df.withColumn("power_lag1_delta", round(col("power_lag1") - col("Total_global_active_power"), 2)) # Create window average fields def add_window_avg_fields(df, window_sizes): for idx, window_size in enumerate(window_sizes, start=1): window_col_name = f"avg_power_lag_{idx}" windowSpec = Window.orderBy("Date").rowsBetween(-window_size, 0) df = df.withColumn(window_col_name, round(avg(col("Total_global_active_power")).over(windowSpec), 2)) return df window_sizes = [14, 30] df = add_window_avg_fields(df, window_sizes) # Create Exponentially Weighted Moving Average (EWMA) fields import pyspark.pandas as ps ps.set_option('compute.ops_on_diff_frames', True) def add_ewma_fields(df, alphas): for idx, alpha in enumerate(alphas, start=1): ewma_col_name = f"ewma_power_weight_{idx}" windowSpec = Window.orderBy("Date") df[ewma_col_name] = df.Total_global_active_power.ewm(alpha=alpha).mean().round(2) return df alphas = [0.2, 0.8] df_pd = df.pandas_api() df_pd = add_ewma_fields(df_pd, alphas) df = df_pd.to_spark() # Write transformed dataframe to the database table "electric_usage_table" df.write.format("jdbc") .option("url", "jdbc:sqlserver://sql-db-dp.database.windows.net:1433;databaseName=sql-db-dp") .option("dbtable", "dbo.electric_usage_table") .option("user", "") .option("password", "") .mode("overwrite") .save()