{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":84493,"databundleVersionId":9871156,"sourceType":"competition"}],"dockerImageVersionId":30822,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"# This Python 3 environment comes with many helpful analytics libraries installed\n# It is defined by the kaggle/python Docker image: https://github.com/kaggle/docker-python\n# For example, here's several helpful packages to load\n\nimport numpy as np # linear algebra\nimport pandas as pd # data processing, CSV file I/O (e.g. pd.read_csv)\n\n# Input data files are available in the read-only \"../input/\" directory\n# For example, running this (by clicking run or pressing Shift+Enter) will list all files under the input directory\n\nimport os\nfor dirname, _, filenames in os.walk('/kaggle/input'):\n    for filename in filenames:\n        print(os.path.join(dirname, filename))\n\n# You can write up to 20GB to the current directory (/kaggle/working/) that gets preserved as output when you create a version using \"Save & Run All\" \n# You can also write temporary files to /kaggle/temp/, but they won't be saved outside of the current session","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2024-12-19T14:01:33.338162Z","iopub.execute_input":"2024-12-19T14:01:33.338496Z","iopub.status.idle":"2024-12-19T14:01:33.76228Z","shell.execute_reply.started":"2024-12-19T14:01:33.338468Z","shell.execute_reply":"2024-12-19T14:01:33.761212Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"###Import Libraries and Initialize Variables\n\nimport pandas as pd\n\n# Initialize a list to hold samples from each file\nsamples = []\n\n###Step 2: Load Data from Multiple Files\n\nfor i in range(10):\n    # Define the path for each partition file\n    file_path = f\"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/partition_id={i}/part-0.parquet\"\n    \n    # Read the data from the Parquet file into a DataFrame\n    chunk = pd.read_parquet(file_path)\n    \n    # Take a sample of 100,000 rows from the current DataFrame\n    sample_chunk = chunk.sample(n=100000, random_state=42)\n    \n    # Append the sampled data to the list\n    samples.append(sample_chunk)\n\n###Step 3: Combine the Samples into a Single DataFrame\n\n# Combine all sampled data into one DataFrame\nsample_df = pd.concat(samples, ignore_index=True)\n\n# Display the first rows of the combined DataFrame\nsample_df.head()\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-19T14:01:41.07903Z","iopub.execute_input":"2024-12-19T14:01:41.079473Z","iopub.status.idle":"2024-12-19T14:03:15.099053Z","shell.execute_reply.started":"2024-12-19T14:01:41.079433Z","shell.execute_reply":"2024-12-19T14:03:15.097885Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"print(sample_df.info())\nprint(sample_df.memory_usage(deep=True).sum())\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-19T13:22:32.432591Z","iopub.execute_input":"2024-12-19T13:22:32.432993Z","iopub.status.idle":"2024-12-19T13:22:32.585838Z","shell.execute_reply.started":"2024-12-19T13:22:32.432961Z","shell.execute_reply":"2024-12-19T13:22:32.584878Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import numpy as np\nimport pandas as pd\nfrom sklearn.preprocessing import StandardScaler\nfrom sklearn.model_selection import train_test_split\nfrom tensorflow.keras.models import Sequential\nfrom tensorflow.keras.layers import LSTM, Dense, Dropout\n\n# Load sample_df (assuming sample_df is already in memory)\n# For demonstration purposes, we'll use a placeholder for sample_df\n\n# Example of sample_df (replace with actual data)\n# sample_df = pd.DataFrame(...)\n\n# Step 1: Feature Engineering\nprint(\"Creating lag and rolling mean features...\")\ndef create_lag_features(df, target_col, lags=[1, 2, 3], rolling_window=3):\n    for lag in lags:\n        df[f'{target_col}_lag{lag}'] = df[target_col].shift(lag)\n    df[f'{target_col}_rolling_mean'] = df[target_col].rolling(rolling_window).mean()\n    return df.dropna()\n\ntarget_col = 'responder_6'\nsample_df = create_lag_features(sample_df, target_col)\n\n# Step 2: Define Features and Target\nprint(\"Defining features and target...\")\nfeatures = [\n    col for col in sample_df.columns if col.startswith('feature_') or 'lag' in col or 'rolling' in col\n]\nX = sample_df[features].to_numpy()\ny = sample_df[target_col].to_numpy()\n\n# Step 3: Scale Features\nprint(\"Scaling features...\")\nscaler = StandardScaler()\nX_scaled = scaler.fit_transform(X)\n\n# Step 4: Reshape Data for LSTM\nprint(\"Reshaping data for LSTM...\")\ntimesteps = 10  # Number of timesteps in each sequence\nX_lstm = np.array([X_scaled[i:i + timesteps] for i in range(len(X_scaled) - timesteps)])\ny_lstm = y[timesteps:]  # Align target with the sequences\n\n# Step 5: Train-Test Split\nprint(\"Splitting data into train and test sets...\")\nX_train, X_test, y_train, y_test = train_test_split(\n    X_lstm, y_lstm, test_size=0.2, random_state=42, shuffle=False\n)\n\n# Step 6: Define LSTM Model\nprint(\"Defining LSTM model...\")\nmodel = Sequential([\n    LSTM(64, return_sequences=True, input_shape=(timesteps, X_train.shape[2])),\n    Dropout(0.2),\n    LSTM(32, return_sequences=False),\n    Dropout(0.2),\n    Dense(16, activation='relu'),\n    Dense(1)  # Single output for regression\n])\n\nmodel.compile(optimizer='adam', loss='mean_squared_error', metrics=['mse'])\nmodel.summary()\n\n# Step 7: Train the Model\nprint(\"Training the model...\")\nhistory = model.fit(\n    X_train, y_train,\n    validation_split=0.2,\n    epochs=20,\n    batch_size=128,\n    verbose=1\n)\n\n# Step 8: Evaluate the Model\nprint(\"Evaluating the model...\")\ntest_loss, test_mse = model.evaluate(X_test, y_test)\nprint(f\"Test MSE: {test_mse}\")\n\n# Step 9: Make Predictions\nprint(\"Making predictions...\")\ny_pred = model.predict(X_test)\n\n# Step 10: Save the Model\nprint(\"Saving the model...\")\nmodel.save(\"time_series_model.h5\")\n\n# Step 11: Visualize Training History\nimport matplotlib.pyplot as plt\n\nplt.plot(history.history['loss'], label='Train Loss')\nplt.plot(history.history['val_loss'], label='Validation Loss')\nplt.legend()\nplt.title('Loss Curves')\nplt.show()\n\nprint(\"Done!\")\n","metadata":{"trusted":true,"execution":{"execution_failed":"2024-12-19T13:17:56.318Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import pandas as pd\nimport numpy as np\nimport tensorflow as tf\nfrom tensorflow.keras.models import Sequential\nfrom tensorflow.keras.layers import LSTM, Dense, Dropout, Input\nfrom sklearn.model_selection import train_test_split\nfrom sklearn.preprocessing import StandardScaler\nfrom sklearn.metrics import mean_squared_error, r2_score\n\n# Step 1: Preprocess Data\nprint(\"Creating lag and rolling mean features...\")\ndef create_lag_features(df, target_col, lags=[1, 2, 3], rolling_window=3):\n    for lag in lags:\n        df[f'{target_col}_lag{lag}'] = df[target_col].shift(lag)\n    df[f'{target_col}_rolling_mean'] = df[target_col].rolling(rolling_window).mean()\n    return df.dropna()\n\ntarget_col = 'responder_6'\nsample_df = create_lag_features(sample_df, target_col)\n\n# Define features and target\nprint(\"Defining features and target...\")\nfeatures = [col for col in sample_df.columns if col.startswith('feature_') or 'lag' in col or 'rolling' in col]\nX = sample_df[features]\ny = sample_df[target_col]\n\n# Scale features\nprint(\"Scaling features...\")\nscaler = StandardScaler()\nX_scaled = scaler.fit_transform(X)\n\n# Reshape data for LSTM\ntimesteps = 10  # Define timesteps for LSTM input\nprint(\"Reshaping data for LSTM...\")\nX_reshaped = np.array([X_scaled[i:i+timesteps] for i in range(len(X_scaled) - timesteps)])\ny_reshaped = y[timesteps:].values\n\n# Train-test split\nprint(\"Splitting data into train and test sets...\")\nX_train, X_test, y_train, y_test = train_test_split(X_reshaped, y_reshaped, test_size=0.2, random_state=42)\n\n# Step 2: Define LSTM Model\nprint(\"Defining LSTM model...\")\nmodel = Sequential([\n    Input(shape=(timesteps, X_train.shape[2])),\n    LSTM(128, return_sequences=True),\n    Dropout(0.2),\n    LSTM(64),\n    Dropout(0.2),\n    Dense(32, activation='relu'),\n    Dense(1)\n])\n\n# Compile the model\noptimizer = tf.keras.optimizers.Adam(learning_rate=0.001)\nmodel.compile(optimizer=optimizer, loss='mse', metrics=['mse'])\nmodel.summary()\n\n# Step 3: Train the Model\nprint(\"Training the model...\")\ncallbacks = [\n    tf.keras.callbacks.ReduceLROnPlateau(patience=3, factor=0.5, min_lr=1e-6),\n    tf.keras.callbacks.EarlyStopping(patience=5, restore_best_weights=True)\n]\n\nhistory = model.fit(\n    X_train, y_train,\n    validation_data=(X_test, y_test),\n    epochs=50,\n    batch_size=128,\n    callbacks=callbacks,\n    verbose=2\n)\n\n# Step 4: Evaluate the Model\nprint(\"Evaluating the model...\")\ny_pred = model.predict(X_test)\nmse = mean_squared_error(y_test, y_pred)\nr2 = r2_score(y_test, y_pred)\n\nprint(f\"Mean Squared Error (MSE): {mse:.4f}\")\nprint(f\"R² Score: {r2:.4f}\")\n\n# Step 5: Save the Model and Scaler\nprint(\"Saving the model and scaler...\")\nmodel.save(\"lstm_model.h5\")\nnp.save(\"scaler_mean.npy\", scaler.mean_)\nnp.save(\"scaler_scale.npy\", scaler.scale_)\n\n# Step 6: Prepare Inference Code\nprint(\"Inference ready. Use 'lstm_model.h5' and scaler files for deployment.\")\n","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-19T13:22:48.524659Z","iopub.execute_input":"2024-12-19T13:22:48.525036Z","iopub.status.idle":"2024-12-19T13:57:03.438979Z","shell.execute_reply.started":"2024-12-19T13:22:48.525008Z","shell.execute_reply":"2024-12-19T13:57:03.437529Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import numpy as np\nimport pandas as pd\nimport polars as pl\nfrom sklearn.preprocessing import MinMaxScaler\nfrom sklearn.model_selection import train_test_split\nfrom keras.models import Sequential\nfrom keras.layers import LSTM, Dense, Dropout\nfrom keras.callbacks import ReduceLROnPlateau, EarlyStopping\nimport matplotlib.pyplot as plt\n\n# Load the sample_df (replace with actual loading if needed)\n# sample_df should already be prepared and available.\n\n# Step 1: Create lag and rolling mean features\ndef create_lag_features(df, lag=10, rolling=5):\n    for i in range(1, lag + 1):\n        df = df.with_columns(df['feature_0'].shift(i).alias(f'lag_{i}'))\n    df = df.with_columns(df['feature_0'].rolling_mean(window_size=rolling).alias(f'rolling_mean_{rolling}'))\n    df = df.drop_nulls()  # Drop rows with NaN values caused by lagging\n    return df\n\nsample_df = create_lag_features(sample_df)\n\n# Step 2: Define features and target\ntarget = \"target\"  # Replace with actual target column name\nfeature_columns = [col for col in sample_df.columns if col != target and col != 'row_id']\n\n# Step 3: Scale features\nscaler = MinMaxScaler()\nscaled_features = scaler.fit_transform(sample_df.select(feature_columns).to_numpy())\nscaled_target = scaler.fit_transform(sample_df.select(target).to_numpy().reshape(-1, 1))\n\n# Step 4: Reshape data for LSTM\nX = scaled_features\ny = scaled_target\ntime_steps = 10  # Sequence length\nnum_features = X.shape[1]\n\nX_lstm = np.array([X[i - time_steps:i] for i in range(time_steps, len(X))])\ny_lstm = y[time_steps:]\n\n# Step 5: Split data into train and test sets\nX_train, X_test, y_train, y_test = train_test_split(X_lstm, y_lstm, test_size=0.2, random_state=42)\n\n# Step 6: Define LSTM model\nmodel = Sequential([\n    LSTM(128, return_sequences=True, input_shape=(time_steps, num_features)),\n    Dropout(0.2),\n    LSTM(64),\n    Dropout(0.2),\n    Dense(32, activation='relu'),\n    Dense(1)\n])\n\nmodel.compile(optimizer='adam', loss='mse', metrics=['mse'])\n\n# Step 7: Define callbacks\nlr_scheduler = ReduceLROnPlateau(monitor='val_loss', factor=0.5, patience=5, verbose=1)\nearly_stopping = EarlyStopping(monitor='val_loss', patience=10, restore_best_weights=True, verbose=1)\n\n# Step 8: Train the model\nhistory = model.fit(\n    X_train, y_train,\n    validation_data=(X_test, y_test),\n    epochs=50,\n    batch_size=128,\n    callbacks=[lr_scheduler, early_stopping],\n    verbose=2\n)\n\n# Step 9: Plot training history\nplt.plot(history.history['loss'], label='Training Loss')\nplt.plot(history.history['val_loss'], label='Validation Loss')\nplt.legend()\nplt.show()\n\n# Step 10: Save the model and scaler for inference\nmodel.save(\"lstm_model.h5\")\nimport pickle\nwith open(\"scaler.pkl\", \"wb\") as f:\n    pickle.dump(scaler, f)\n\n# Step 11: Define prediction function for inference server\ndef predict(test: pl.DataFrame, lags: pl.DataFrame | None) -> pl.DataFrame | pd.DataFrame:\n    global scaler  # Assuming the scaler is loaded globally\n    global model   # Assuming the model is loaded globally\n\n    # Select feature columns\n    feature_columns = [col for col in test.columns if col.startswith(\"feature_\")]\n    features = test.select(feature_columns).to_numpy()\n    \n    # Scale features\n    features_scaled = scaler.transform(features)\n\n    # Reshape for LSTM\n    time_steps = 10\n    num_features = features_scaled.shape[1]\n    X_lstm_inference = np.array([features_scaled[i - time_steps:i] for i in range(time_steps, len(features_scaled))])\n\n    # Generate predictions\n    predictions = model.predict(X_lstm_inference)\n    \n    # Construct output DataFrame\n    output_df = test.select(\"row_id\").slice(time_steps).with_columns(\n        pl.Series(\"target\", predictions.flatten())\n    )\n    \n    return output_df\n\n# Step 12: Ready for inference server integration\nimport os\nimport kaggle_evaluation.jane_street_inference_server\n\ninference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\nif os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n    inference_server.serve()\nelse:\n    inference_server.run_local_gateway(\n        (\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,"execution":{"iopub.status.busy":"2024-12-19T14:03:36.920485Z","iopub.execute_input":"2024-12-19T14:03:36.920883Z","iopub.status.idle":"2024-12-19T14:03:49.039517Z","shell.execute_reply.started":"2024-12-19T14:03:36.920847Z","shell.execute_reply":"2024-12-19T14:03:49.03813Z"}},"outputs":[],"execution_count":null}]}