{"metadata":{"kernelspec":{"name":"ir","display_name":"R","language":"R"},"language_info":{"name":"R","codemirror_mode":"r","pygments_lexer":"r","mimetype":"text/x-r-source","file_extension":".r","version":"4.4.0"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30749,"isInternetEnabled":false,"language":"r","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# Much of this notebook was adapted from a Python public notebook to be referenced soon\nlibrary(arrow)\nlibrary(tidyverse)\nlibrary(catboost)\nlibrary(R6)\nlibrary(reticulate)\n# Suppress warnings\noptions(warn = -1)","metadata":{"_uuid":"051d70d956493feee0c6d64651c6a088724dca2a","_execution_state":"idle","execution":{"iopub.status.busy":"2024-11-13T13:50:14.028398Z","iopub.execute_input":"2024-11-13T13:50:14.030124Z","iopub.status.idle":"2024-11-13T13:50:17.030592Z","shell.execute_reply":"2024-11-13T13:50:17.028968Z"},"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Use reticulate to run pip directly for the .whl file\nreticulate::py_run_string(\"\nimport pip\npip.main(['install', '/kaggle/input/polarspackage/polars-1.12.0-cp39-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl'])\n\")","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-13T13:50:17.03305Z","iopub.execute_input":"2024-11-13T13:50:17.079358Z","iopub.status.idle":"2024-11-13T13:50:39.691352Z","shell.execute_reply":"2024-11-13T13:50:39.689662Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"os = import(\"os\")\npd = import(\"pandas\")\n\nimport_directory = \"/kaggle/input/jane-street-real-time-market-data-forecasting\"","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-11-13T13:50:39.693845Z","iopub.execute_input":"2024-11-13T13:50:39.69514Z","iopub.status.idle":"2024-11-13T13:50:40.846852Z","shell.execute_reply":"2024-11-13T13:50:40.831392Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"lags_r <- arrow::read_parquet(\"/kaggle/input/jane-street-real-time-market-data-forecasting/lags.parquet/date_id=0/part-0.parquet\")\nlags <- pl$DataFrame(r_to_py(lags_r))  # Convert to polars DataFrame\n\ntest_r <- arrow::read_parquet(\"/kaggle/input/jane-street-real-time-market-data-forecasting/test.parquet/date_id=0/part-0.parquet\")\ntest <- pl$DataFrame(r_to_py(test_r))  # Convert to polars DataFrame","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Define the Python function using reticulate's `py_run_string`\npy_run_string(\"\nimport polars as pl\n\nlags_ = None\n\ndef predict(test, lags):\n    global lags_\n    if lags is not None:\n        lags_ = lags\n\n    # Replace this section with your own predictions\n    predictions = test.select(\n        'row_id',\n        pl.lit(0.0).alias('responder_6'),\n    )\n\n    if isinstance(predictions, pl.DataFrame):\n        assert predictions.columns == ['row_id', 'responder_6']\n    elif isinstance(predictions, pd.DataFrame):\n        assert (predictions.columns == ['row_id', 'responder_6']).all()\n    else:\n        raise TypeError('The predict function must return a DataFrame')\n    # Confirm it has as many rows as the test data\n    assert len(predictions) == len(test)\n\n    return predictions\n\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Load and prepare training data\ntrain <- read_parquet(\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/partition_id=9/part-0.parquet\")\n\n# Calculate features using all columns from 4 to the end\nfeatures_agg_sum <- train %>%\n  select(4:ncol(train)) %>%\n  colSums() %>%\n  enframe(name = \"feature\", value = \"Value\")\n\nfeatures_agg_sum_large <- features_agg_sum %>%\n  filter(Value > 2)\n\nfeature_names <- features_agg_sum_large$feature\n\n# Create ModelSelector class using R6\nModelSelector <- R6Class(\n  \"ModelSelector\",\n  public = list(\n    threshold = NULL,\n    feature_importance = NULL,\n    best_model = NULL,\n    \n    initialize = function(threshold = 0.8) {\n      self$threshold <- threshold\n    },\n    \n    train_and_evaluate = function(train_data, train_y, feature_names) {\n      # Create validation split (20% of data)\n      val_size <- floor(nrow(train_data) * 0.2)\n      train_idx <- nrow(train_data) - val_size\n      \n      train_subset <- train_data[1:train_idx, ]\n      train_y_subset <- train_y[1:train_idx]\n      val_subset <- train_data[(train_idx + 1):nrow(train_data), ]\n      val_y_subset <- train_y[(train_idx + 1):length(train_y)]\n      \n      # Create proper pool objects for CatBoost\n      train_pool <- catboost.load_pool(\n        data = as.matrix(train_subset[, feature_names, drop = FALSE]),\n        label = as.numeric(train_y_subset)\n      )\n      \n      val_pool <- catboost.load_pool(\n        data = as.matrix(val_subset[, feature_names, drop = FALSE]),\n        label = as.numeric(val_y_subset)\n      )\n      \n      # Initialize base model with proper parameters\n      params <- list(\n        loss_function = \"RMSE\",\n        iterations = 3000,\n        learning_rate = 0.03,\n        depth = 8,\n        l2_leaf_reg = 3,\n        min_data_in_leaf = 50,\n        eval_metric = \"RMSE\",\n        verbose = 100\n      )\n      \n      # Train base model\n      base_model <- catboost.train(\n        learn_pool = train_pool,\n        test_pool = val_pool,\n        params = params\n      )\n      \n      # Calculate feature importance\n      importance <- catboost.get_feature_importance(base_model, pool = train_pool)\n      self$feature_importance <- data.frame(\n        feature = feature_names,\n        importance = importance\n      ) %>%\n        arrange(desc(importance))\n      \n      # Calculate R² score on validation set\n      pred <- catboost.predict(base_model, val_pool)\n      r2 <- self$calculate_weighted_r2(val_y_subset, pred)\n      \n      # If R² is below threshold, create ensemble\n      if (r2 < self$threshold) {\n        cat(sprintf(\"R² score (%.4f) below threshold (%.1f). Creating ensemble...\\n\", \n                   r2, self$threshold))\n        self$best_model <- self$create_ensemble(train_data, train_y, feature_names)\n      } else {\n        cat(sprintf(\"R² score (%.4f) above threshold. Using single model.\\n\", r2))\n        self$best_model <- base_model\n      }\n    },\n    \n    calculate_weighted_r2 = function(y_true, y_pred, weights = NULL) {\n      if (is.null(weights)) weights <- rep(1, length(y_true))\n      numerator <- sum(weights * (y_true - y_pred)^2)\n      denominator <- sum(weights * y_true^2)\n      return(1 - (numerator / denominator))\n    },\n    \n    create_ensemble = function(train_data, train_y, feature_names) {\n      # Base parameters for all models\n      params <- list(\n        loss_function = \"RMSE\",\n        iterations = 300,\n        learning_rate = 0.05,\n        depth = 6,\n        l2_leaf_reg = 5,\n        bagging_temperature = 1,\n        border_count = 128,\n        grow_policy = 'SymmetricTree',\n        use_best_model = TRUE,\n        min_data_in_leaf = 50,\n        eval_metric = \"RMSE\",\n        verbose = 100\n      )\n      \n      # Train multiple models with different seeds\n      models <- list()\n      seeds <- c(42, 84, 168)\n      \n      train_pool <- catboost.load_pool(\n        data = as.matrix(train_data[, feature_names, drop = FALSE]),\n        label = as.numeric(train_y)\n      )\n      \n      for (seed in seeds) {\n        params$random_seed <- seed\n        model <- catboost.train(\n          learn_pool = train_pool,\n          params = params\n        )\n        models[[length(models) + 1]] <- model\n      }\n      \n      return(models)\n    },\n    \n    predict = function(test_data, feature_names) {\n      test_pool <- catboost.load_pool(\n        data = as.matrix(test_data[, feature_names, drop = FALSE])\n      )\n      \n      if (is.list(self$best_model)) {\n        # For ensemble, average predictions\n        predictions <- sapply(self$best_model, function(model) {\n          catboost.predict(model, test_pool)\n        })\n        return(rowMeans(predictions))\n      } else {\n        # For single model\n        return(catboost.predict(self$best_model, test_pool))\n      }\n    },\n    \n    print_feature_importance = function() {\n      if (!is.null(self$feature_importance)) {\n        cat(\"\\nTop 10 Most Important Features:\\n\")\n        print(head(self$feature_importance, 10))\n      } else {\n        cat(\"Feature importance not available. Train the model first.\\n\")\n      }\n    }\n  )\n)\n\n# Initialize ModelSelector\nmodel_selector <- ModelSelector$new(threshold = 0.8)\n\n# Prepare target variable\ntrain_y <- train$responder_6\n\n# Train the model using only selected features\nmodel_selector$train_and_evaluate(\n  train_data = train,\n  train_y = train_y,\n  feature_names = feature_names\n)\n\n# Print feature importance after training\nmodel_selector$print_feature_importance()\n\n# Global variable for lags\nlags_ <- NULL\n\n# Prediction function\npredict <- function(test, lags = NULL) {\n  if (!is.null(lags)) {\n    lags_ <<- lags\n  }\n  \n  # Create predictions dataframe\n  predictions <- tibble(\n    row_id = test$row_id,\n    responder_6 = 0.0\n  )\n  \n  # Make predictions using model\n  pred <- model_selector$predict(test, feature_names)\n  predictions$responder_6 <- pred\n  \n  # Assertions for validation\n  stopifnot(\n    is.data.frame(predictions),\n    identical(names(predictions), c(\"row_id\", \"responder_6\")),\n    nrow(predictions) == nrow(test)\n  )\n  \n  return(predictions)\n}","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Create a Python function that wraps the R predict function\npy_run_string(\"\ndef predict_wrapper(market_obs_df):\n    import pandas as pd\n    from rpy2.robjects import pandas2ri\n    pandas2ri.activate()\n    \n    # Convert Python DataFrame to R DataFrame\n    r_df = pandas2ri.py2rpy(market_obs_df)\n    \n    # Call R predict function\n    result = r.predict(r_df)\n    \n    return float(result)\n\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"kaggle_evaluation <- import_from_path(\"kaggle_evaluation.jane_street_inference_server\", path = import_directory)\n# Define the inference server with the predict function\ninference_server <- kaggle_evaluation$JSInferenceServer(py$predict)\n\n# Check the environment variable and decide which method to call\nif (Sys.getenv(\"KAGGLE_IS_COMPETITION_RERUN\") != \"\") {\n    inference_server$serve()\n} else {\n    inference_server$run_local_gateway(\n        list(\n            \"/kaggle/input/jane-street-real-time-market-data-forecasting/test.parquet\",\n            \"/kaggle/input/jane-street-real-time-market-data-forecasting/lags.parquet\"\n        )\n    )\n}","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}