{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.10.14","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"},{"sourceId":10044851,"sourceType":"datasetVersion","datasetId":6182184}],"dockerImageVersionId":30786,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import pandas as pd\nimport numpy as np\nimport lightgbm as lgb\nimport time\nimport polars as pl\nimport os\nimport shutil\nfrom tqdm import tqdm\nimport kaggle_evaluation.jane_street_inference_server\n\ndo_predict=True\n\ninput_dir = '/kaggle/input/jane-street-real-time-market-data-forecasting/'\nmodel_save_path = '/kaggle/working/'\n#input_dir = '/data/release/kaggle-jane/'\n#model_save_path = 'model_save/'\ntrain_path = input_dir + 'train.parquet'\ntest_path = input_dir + 'test.parquet'\nlag_path = input_dir + 'lags.parquet'\ntrain_data_files = None\ntest_data_files = None\nlags_ = None\ntrain_kfold = False\nfold_num = 5\n\nfeature_list = [f\"feature_{idx:02d}\" for idx in range(79)]\nselected_features = ['symbol_id', 'time_id'] + [f\"responder_{idx}_lag_1\" for idx in range(9)] + feature_list \ntarget_col = 'responder_6'\n\nmodel_dir = '/kaggle/input/jane-street-lgb-model-used-data-6-to-10/'\nmodel_name = 'lgb_model'\n\nif os.path.exists(model_save_path):\n    shutil.rmtree(model_save_path + '/*', ignore_errors=True)\n    print('Path %s was cleared.'%model_save_path)\n    \n\ndef read_data():\n    # to run on Kaggle, the maximun data range is (2,10), otherwise it would run out of the memory\n    file_paths = [train_path + '/partition_id=' + str(i) + '/part-0.parquet' for i in range(2,10)]\n    train_df = pl.concat([pl.scan_parquet(fp) for fp in file_paths])\n    return train_df \n\ndef create_lag_features(df):\n    df = df.sort(['date_id', 'time_id', 'symbol_id'])\n\n    for idx in range(9):\n        original_col = f\"responder_{idx}\"\n        new_col = f\"responder_{idx}_lag_1\"\n        df = df.with_columns(\n            pl.col(original_col)\n            .shift(1)\n            .over(['time_id', 'symbol_id'])\n            .alias(new_col)\n        )\n    return df\n\ndef create_lag_features2(df):\n    lag_cols_original = [\"date_id\", \"symbol_id\"] + [f\"responder_{idx}\" for idx in range(9)]\n    lag_cols_rename = { f\"responder_{idx}\" : f\"responder_{idx}_lag_1\" for idx in range(9)}\n\n    lags = df.select(pl.col(lag_cols_original))\n    lags = lags.rename(lag_cols_rename)\n    lags = lags.with_columns(\n        date_id = pl.col('date_id') + 1,  # lagged by 1 day\n        )\n    lags = lags.group_by([\"date_id\", \"symbol_id\"], maintain_order=True).last()\n    df = df.join(lags, on=[\"date_id\", \"symbol_id\"],  how=\"left\")\n    return df\n\ndef weighted_zero_mean_r2(y_true, y_pred, weights):\n    numerator = np.sum(weights * (y_true - y_pred)**2)\n    denominator = np.sum(weights * y_true**2)\n    \n    r2_score = 1 - numerator / denominator\n    return r2_score\n\ndef batch_load_data(df, start_date, end_date, selected_features, target_col, batch_size, output_dir):\n    if not os.path.exists(output_dir):\n        os.makedirs(output_dir)\n    \n    file_list = []\n    for batch_start in tqdm(range(start_date, end_date + 1, batch_size), desc=\"Loading data batches\"):\n        batch_end = min(batch_start + batch_size - 1, end_date)\n        batch_df = df.filter((pl.col('date_id') >= batch_start) & (pl.col('date_id') <= batch_end))\n        batch_data = batch_df.select(selected_features + [target_col, 'weight']).collect().to_numpy()\n        \n        file_path = os.path.join(output_dir, f'batch_{batch_start}_{batch_end}.npz')\n        np.savez(file_path, data=batch_data)\n        file_list.append(file_path)\n        \n        del batch_df, batch_data\n    return file_list\n    \ndef split_data(df, test_ratio=0.2, batch_size=100, output_dir='data_batches'):\n    max_date_id = df.select(pl.col('date_id').max()).collect()[0, 0]\n    min_date_id = df.select(pl.col('date_id').min()).collect()[0, 0]\n    date_split = max_date_id - int((max_date_id - min_date_id) * test_ratio)\n    \n    print('Train data from [%d, %d]' % (min_date_id, date_split - 1))\n    print('Test data from [%d, %d]' % (date_split, max_date_id))\n    \n    train_data_files = batch_load_data(df, min_date_id, date_split - 1, selected_features, target_col, batch_size, output_dir)\n    test_data_files = batch_load_data(df, date_split, max_date_id, selected_features, target_col, batch_size, output_dir)\n    \n    return train_data_files, test_data_files\n\ndef get_kfold_data(df, fold_idx, selected_features, target_col, batch_size=100, output_dir='kfold_data_batch', fold_num=5, gap=1,):\n    max_date_id = df.select(pl.col('date_id').max()).collect()[0, 0]\n    min_date_id = df.select(pl.col('date_id').min()).collect()[0, 0]\n    fold_size = (max_date_id  - min_date_id + 1) // fold_num\n    if fold_idx == 0:\n        print('Data date id range from %d to %d. Fold size: %d.'%(min_date_id, max_date_id, fold_size))\n\n    if fold_idx < fold_num - 1:\n        start = fold_idx * fold_size\n        end = start + fold_size\n        purged_start = end - gap - 1\n        purged_end = end + fold_size + gap\n        if purged_end > max_date_id:\n            purged_end = max_date_id\n        last_idx = min(end + fold_size, max_date_id) - 1\n        print('fold num: %d, 1st train set: [%s, %s], 2nd train set: [%s, %s], test set: [%s, %s]'\n                %(fold_idx, start, purged_start, purged_end, max_date_id, end, last_idx))\n        train_data_files_1st = batch_load_data(df, start, purged_start, selected_features, target_col, batch_size, output_dir)\n        tarin_data_files_2nd = batch_load_data(df, purged_end, max_date_id, selected_features, target_col, batch_size, output_dir)\n        train_data_files = train_data_files_1st + tarin_data_files_2nd\n        train_data = load_data_from_files(train_data_files)\n\n        test_data_files = batch_load_data(df, end, last_idx, selected_features, target_col, batch_size, output_dir)\n        test_data = load_data_from_files(test_data_files)\n    else:\n        train_day_start = fold_size + gap\n        test_day_end = fold_size - gap\n        print('fold num: %d, train set: [%s, %s], test set: [%s, %s]'%(fold_num - 1, train_day_start, max_date_id, min_date_id, test_day_end))\n        train_data_files = batch_load_data(df, train_day_start, max_date_id, selected_features, target_col, batch_size, output_dir)\n        train_data = load_data_from_files(train_data_files)\n        test_data_file = batch_load_data(df, min_date_id, test_day_end, selected_features, target_col, batch_size, output_dir)\n        test_data = load_data_from_files(test_data_file)\n    return train_data, test_data\n\ndef load_data_from_files(file_list):\n    if not file_list:\n        return np.array([])\n    \n    with np.load(file_list[0]) as npz_file:\n        data = npz_file['data']\n    \n    for file_path in file_list[1:]:\n        with np.load(file_path) as npz_file:\n            batch_data = npz_file['data']\n            data = np.vstack((data, batch_data))\n    return data\n\ndef load_train_test_data(train_data_files, test_data_files):\n    print('Load train files: %s'%train_data_files)\n    train_data = load_data_from_files(train_data_files)\n    print('Load test files: %s'%test_data_files)\n    test_data = load_data_from_files(test_data_files)\n    print('train len: %d, test len: %d' % (len(train_data), len(test_data)))\n    return train_data, test_data\n\ndef train_lgb_model(train_data, test_data, idx=0):\n    params = {\n        \"objective\": \"mae\",\n        \"n_estimators\": 6000,\n        \"num_leaves\": 256,\n        \"subsample\": 0.6,\n        \"colsample_bytree\": 0.8,\n        \"learning_rate\": 0.01,\n        'max_depth': 11,\n        \"n_jobs\": -1,\n        \"verbosity\": -1,\n        \"importance_type\": \"gain\",\n        \"reg_alpha\": 0.2,\n        \"reg_lambda\": 3.25\n    }\n    x = train_data[:, :-2]\n    y = train_data[:, -2]\n    test_x = test_data[:, :-2]\n    test_y = test_data[:, -2]\n    test_weight = test_data[:, -1]\n    \n    model = lgb.LGBMRegressor(**params)\n    t1 = time.time()\n    bst = model.fit(x, y,\n        eval_set=[(test_x, test_y)],\n        callbacks=[\n            lgb.callback.early_stopping(stopping_rounds=200),\n            lgb.callback.log_evaluation(period=200),\n        ],\n    )\n    t2 = time.time()\n    print(\"Training elapsed time:%d sec.\" % (t2 - t1))\n    model.booster_.save_model(model_save_path + 'lgb_model_' + str(idx))\n    return model\n\ndef train_single_lgb_model():\n    global train_data_files, test_data_files\n    if train_data_files is None or test_data_files is None:\n        df = read_data()\n        df = create_lag_features(df)\n        train_data_files, test_data_files = split_data(df)\n        del df\n    train_data, test_data = load_train_test_data(train_data_files, test_data_files)\n    model = train_lgb_model(train_data, test_data)\n    test_x = test_data[:, :-2]\n    test_y = test_data[:, -2]\n    test_weight = test_data[:, -1]\n    test_pred = model.predict(test_x)\n    print(weighted_zero_mean_r2(test_y, test_pred, test_weight))\n    return model\n\ndef train_kfold_lgb_model(fold_num=5):\n    df = read_data()\n    df = create_lag_features(df)\n\n    models = []\n    for i in range(fold_num):\n        train_data, test_data = get_kfold_data(df, i, selected_features, target_col, fold_num=fold_num)\n        model = train_lgb_model(train_data, test_data, i)\n        models.append(model)\n    return models\n\ndef load_model():\n    if train_kfold:\n        model = []\n        for i in range(fold_num):\n            load_path = model_dir + model_name + '_' + str(i)\n            model.append(lgb.Booster(model_file=load_path))\n    else:\n        load_path = model_dir + model_name + '_0'\n        model = lgb.Booster(model_file=load_path)\n    print('Load model from the path:%s'%model_dir)\n    return model\n\ndef predict(test: pl.DataFrame, lags: pl.DataFrame):\n    global lags_  \n    \n    if lags is not None:\n        lags_ = lags\n\n    predictions = test.select(\n        'row_id',  \n        pl.lit(0.0).alias('responder_6'),  \n    )\n    if lags_ is not None:\n        test = test.join(lags_, on=[\"date_id\", \"time_id\", \"symbol_id\"], how=\"left\") \n    else:\n        # If no lag data is provided, create columns with default 0 values for lag features\n        test = test.with_columns(\n            (pl.lit(0.0).alias(f'responder_{idx}_lag_1') for idx in range(9)) \n        )\n   \n    # Initialize a zero array for predictions\n    preds = np.zeros((test.shape[0],))\n    if train_kfold:\n        for m in model:\n            preds += m.predict(test.select(selected_features).to_pandas().values)\n        preds /= len(model)\n    else:\n        preds = model.predict(test.select(selected_features).to_pandas().values)\n\n    # Generate the final prediction DataFrame, ensuring predictions are clipped between -5 and 5\n    predictions = test.select('row_id').with_columns(\n        pl.Series(\n            name='responder_6',  \n            values=np.clip(preds, a_min=-5, a_max=5),  \n            dtype=pl.Float64,  \n        )\n    )\n\n    # Ensure the prediction function returns a DataFrame\n    assert isinstance(predictions, pl.DataFrame | pd.DataFrame)\n\n    # Ensure the returned DataFrame has columns 'row_id' and 'responder_6'\n    assert list(predictions.columns) == ['row_id', 'responder_6']\n\n    # Ensure the number of rows in the prediction matches the test data\n    assert len(predictions) == len(test)\n\n    return predictions \n\n\nif do_predict:\n    model = load_model()\nelse:\n    if train_kfold:\n        model = train_kfold_lgb_model(fold_num)\n    else:\n        model = train_single_lgb_model()\n\n# predict\ninference_server = kaggle_evaluation.jane_street_inference_server.JSInferenceServer(predict)\nif os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n    print('Online predict...')\n    inference_server.serve() \nelse:\n    print('Local predict...')\n    # If not in the competition rerun environment, run a local gateway with the provided test and lag data\n    inference_server.run_local_gateway(\n        (\n            test_path,\n            lag_path, \n        )\n    )\nprint('All done.')","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2024-11-29T02:33:06.143071Z","iopub.execute_input":"2024-11-29T02:33:06.143691Z","iopub.status.idle":"2024-11-29T02:33:06.470350Z","shell.execute_reply.started":"2024-11-29T02:33:06.143648Z","shell.execute_reply":"2024-11-29T02:33:06.469146Z"}},"outputs":[],"execution_count":null}]}