{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":56537,"databundleVersionId":8015876,"sourceType":"competition"}],"dockerImageVersionId":30698,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"pip install pyspark ","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:11:22.18891Z","iopub.execute_input":"2024-04-21T04:11:22.189306Z","iopub.status.idle":"2024-04-21T04:11:37.020442Z","shell.execute_reply.started":"2024-04-21T04:11:22.189273Z","shell.execute_reply":"2024-04-21T04:11:37.018928Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import numpy as np\nimport pandas as pd\nimport time\n\n# plots\nimport matplotlib.pyplot as plt\nimport seaborn as sns","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","execution":{"iopub.status.busy":"2024-04-21T04:11:37.023623Z","iopub.execute_input":"2024-04-21T04:11:37.024161Z","iopub.status.idle":"2024-04-21T04:11:38.069969Z","shell.execute_reply.started":"2024-04-21T04:11:37.024094Z","shell.execute_reply":"2024-04-21T04:11:38.06868Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# configs\npd.set_option('display.max_columns', None) # we want to display all columns in this notebook\n\n# aesthetics\ndefault_color_1 = 'darkblue'\ndefault_color_2 = 'darkgreen'\ndefault_color_3 = 'darkred'\n\n# random seed\nmy_random_seed = 111","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:11:38.071445Z","iopub.execute_input":"2024-04-21T04:11:38.071984Z","iopub.status.idle":"2024-04-21T04:11:38.078958Z","shell.execute_reply.started":"2024-04-21T04:11:38.071951Z","shell.execute_reply":"2024-04-21T04:11:38.07768Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"!ls -l '../input/leap-atmospheric-physics-ai-climsim/'\n","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:11:38.08171Z","iopub.execute_input":"2024-04-21T04:11:38.08217Z","iopub.status.idle":"2024-04-21T04:11:39.230703Z","shell.execute_reply.started":"2024-04-21T04:11:38.082106Z","shell.execute_reply":"2024-04-21T04:11:39.229236Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.sql import SparkSession\n\n# Create SparkSession\nspark = SparkSession.builder \\\n    .appName(\"Subset Data Reader\") \\\n    .getOrCreate()\n\n# Set number of rows to read\nn_rows = 5000\n\n# Path to the data folder\nfolder = 'leap-atmospheric-physics-ai-climsim'\n\n# Read data into Spark DataFrame\ndf_train_spark = spark.read.option(\"header\", \"true\").csv(f\"../input/{folder}/train.csv\").limit(n_rows)\ndf_test_spark = spark.read.option(\"header\", \"true\").csv(f\"../input/{folder}/test.csv\")\ndf_sub_spark = spark.read.option(\"header\", \"true\").csv(f\"../input/{folder}/sample_submission.csv\").limit(n_rows)\n\n","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:11:39.232528Z","iopub.execute_input":"2024-04-21T04:11:39.232919Z","iopub.status.idle":"2024-04-21T04:11:53.174826Z","shell.execute_reply.started":"2024-04-21T04:11:39.232883Z","shell.execute_reply":"2024-04-21T04:11:53.173154Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"df_test_spark\ndf_train_spark","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:12:13.268494Z","iopub.execute_input":"2024-04-21T04:12:13.268945Z","iopub.status.idle":"2024-04-21T04:12:13.342994Z","shell.execute_reply.started":"2024-04-21T04:12:13.268903Z","shell.execute_reply":"2024-04-21T04:12:13.341559Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.sql.functions import col\n\n# Assuming df_test_spark and df_train_spark are your Spark DataFrames\n\n# Cast every column in df_test_spark to float\ndf_test_spark = df_test_spark.select(\n    *[col(column).cast(\"float\").alias(column) for column in df_test_spark.columns]\n)\n\n# Cast every column in df_train_spark to float\ndf_train_spark = df_train_spark.select(\n    *[col(column).cast(\"float\").alias(column) for column in df_train_spark.columns]\n)\n","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:13:01.05356Z","iopub.execute_input":"2024-04-21T04:13:01.054627Z","iopub.status.idle":"2024-04-21T04:13:05.547088Z","shell.execute_reply.started":"2024-04-21T04:13:01.054583Z","shell.execute_reply":"2024-04-21T04:13:05.545876Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.sql.functions import col\n\n# Targets (extract from submission file)\ntargets = df_sub_spark.columns[1:]  # Assuming the first column is 'sample_id'\n\n# Numerical features\nfeature_num = [col_name for col_name in df_train_spark.columns if col_name not in ['sample_id'] + targets]\n\n# Categorical features (assuming there are no categorical features for now)\nfeature_cat = []\n\n# All features combined\nfeatures = feature_num + feature_cat\n","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:13:26.522061Z","iopub.execute_input":"2024-04-21T04:13:26.522488Z","iopub.status.idle":"2024-04-21T04:13:26.582722Z","shell.execute_reply.started":"2024-04-21T04:13:26.522457Z","shell.execute_reply":"2024-04-21T04:13:26.581625Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#we need to get variation of the all columns since our sample size cant be capturing full data \n#we are going to build models for sets of datasets the rational i am trying to make is each \n#model will capture the different data pattern and at the end of whole data we will have models \n#trained on specific data\n#lets simulate this idea with simple Linear Regression\n","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from pyspark.ml.feature import VectorAssembler\nfrom pyspark.ml.regression import GBTRegressor\nfrom pyspark.ml.evaluation import RegressionEvaluator\nfrom pyspark.sql.functions import col\n\nclass Modelling:\n    def __init__(self, spark, df_train, df_test, modelling_features, targets):\n        self.spark = spark\n        self.df_train = df_train\n        self.df_test = df_test\n        self.modelling_features = modelling_features\n        self.targets = targets\n\n    def baseline_modelling(self):\n        # Cast columns to numeric type\n        for col_name in self.modelling_features:\n            self.df_train = self.df_train.withColumn(col_name, col(col_name).cast(\"float\"))\n            self.df_test = self.df_test.withColumn(col_name, col(col_name).cast(\"float\"))\n        \n        # Assemble features\n        assembler = VectorAssembler(inputCols=self.modelling_features, outputCol=\"features\")\n        \n        # Train models for each target\n        self.models = {}\n        for target_column in self.targets:\n            train_data = assembler.transform(self.df_train).select(\"features\", target_column)\n            gbt = GBTRegressor(featuresCol=\"features\", labelCol=target_column, maxIter=10)\n            model = gbt.fit(train_data)\n            self.models[target_column] = model\n\n    def prediction_model(self, output_path):\n        # Assemble features for test data\n        assembler = VectorAssembler(inputCols=self.modelling_features, outputCol=\"features\")\n        test_data = assembler.transform(self.df_test).select(\"sample_id\", \"features\")\n\n        # Make predictions for each target\n        predictions = None\n        for target_column, model in self.models.items():\n            result = model.transform(test_data).select(\"sample_id\", \"prediction\").withColumnRenamed(\"prediction\", target_column)\n            if predictions is None:\n                predictions = result\n            else:\n                predictions = predictions.join(result, \"sample_id\")\n\n        # Write predictions to a single CSV file\n        predictions.coalesce(1).write.csv(output_path, header=True)\n\n    def eval(self, predictions):\n        evaluator = RegressionEvaluator(labelCol=\"label\", predictionCol=\"prediction\", metricName=\"r2\")\n        r2_scores = {}\n        for target_column in self.targets:\n            r2_score = evaluator.evaluate(predictions, {evaluator.metricName: target_column})\n            r2_scores[target_column] = r2_score\n        mean_r2 = sum(r2_scores.values()) / len(r2_scores)\n        return mean_r2\n","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:13:29.678029Z","iopub.execute_input":"2024-04-21T04:13:29.678478Z","iopub.status.idle":"2024-04-21T04:13:29.871092Z","shell.execute_reply.started":"2024-04-21T04:13:29.678444Z","shell.execute_reply":"2024-04-21T04:13:29.869913Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Assuming you have defined the Modelling class with Spark support\n\n# 1. Create an instance of Modelling class\nmodel = Modelling(spark, df_train_spark, df_test_spark, feature_num, targets)\n\n# 2. Train the models on df_train\nmodel.baseline_modelling()\n\n# 3. Make predictions on df_test\npredictions = model.prediction_model(\"predictions.csv\")\n\n# 4. Evaluate the predictions\nmean_r2 = model.eval(predictions)\nprint(\"Mean R^2 score:\", mean_r2)\n","metadata":{"execution":{"iopub.status.busy":"2024-04-21T04:13:36.815856Z","iopub.execute_input":"2024-04-21T04:13:36.816437Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"","metadata":{"trusted":true},"execution_count":null,"outputs":[]}]}