spark = SparkSession.builder.appName("CaloriesAnalysis").getOrCreate()def data_cleaning_and_stats_analysis(df): df = df.na.drop(subset=["Calories", "Duration", "Heart_Rate", "Body_Temp"]) df = df.filter((col("Duration") > 0) & (col("Calories") > 0)) stats = df.select( mean("Calories").alias("avg_calories"), stddev("Calories").alias("stddev_calories"), mean("Duration").alias("avg_duration"), mean("Heart_Rate").alias("avg_heart_rate") ).collect()[0] duration_bins = [0, 10, 30, 60, float('inf')] bucketizer = Bucketizer(splits=duration_bins, inputCol="Duration", outputCol="duration_group") df_binned = bucketizer.transform(df) distribution = df_binned.groupBy("duration_group").count().orderBy("duration_group").collect() return {"stats": stats, "distribution": distribution}def correlation_analysis(df): assembler = VectorAssembler( inputCols=["Age", "Weight", "Duration", "Heart_Rate", "Body_Temp"], outputCol="features" ) vec_df = assembler.transform(df).select("features") correlation_matrix = Correlation.corr(vec_df, "features").head()[0].toArray() corr_list = [] col_names = ["Age", "Weight", "Duration", "Heart_Rate", "Body_Temp"] for i in range(len(col_names)): for j in range(len(col_names)): corr_list.append({ "col1": col_names[i], "col2": col_names[j], "corr": float(correlation_matrix[i, j]) }) return corr_listdef kmeans_clustering_analysis(df): assembler = VectorAssembler( inputCols=["Duration", "Heart_Rate", "Calories"], outputCol="features" ) data = assembler.transform(df).select("features") kmeans = KMeans(k=3, seed=1, maxIter=20) model = kmeans.fit(data) predictions = model.transform(data) centers = model.clusterCenters() cluster_stats = predictions.groupBy("prediction").agg( mean("Calories").alias("avg_cal"), mean("Duration").alias("avg_dur"), count("*").alias("count") ).collect() result = [] for row in cluster_stats: result.append({ "cluster_id": row["prediction"], "avg_calories": row["avg_cal"], "avg_duration": row["avg_dur"], "user_count": row["count"] }) return result