{"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"}],"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"#imports and options setting\nimport pandas as pd\n#import seaborn as sns\nimport matplotlib.pyplot as plt\nimport numpy as np\n#import statsmodels.api as sm\nimport polars as pl\nimport polars.selectors as cs\nimport re\nfrom enum import Enum\npd.set_option('display.max_rows', 500)\npd.set_option('display.max_columns', 500)\npd.set_option('display.width', 150)\n\nclass VarCategory(Enum):\n    #Useful enum for later\n    RESPONDER=1\n    FEATURE=2\n#!pip install polars==1.16.0\n#print(pl.show_versions())","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2024-12-17T16:22:01.032659Z","iopub.execute_input":"2024-12-17T16:22:01.033090Z","iopub.status.idle":"2024-12-17T16:22:03.048756Z","shell.execute_reply.started":"2024-12-17T16:22:01.033052Z","shell.execute_reply":"2024-12-17T16:22:03.047622Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"#Building functions for manipulating responders & features in a pipeline\ndef responder_level_flags(dataframe,level):\n    #Flagging levels in responder variables (already known that the responders are capped between -5 and 5)\n    dataframe=dataframe.with_columns((abs(pl.col(\"^responder_\\d?$\")) >= level).name.suffix(\"_gt\"+str(level)+\"_abs\"),\n                           (pl.col(\"^responder_\\d?$\") >= level).name.suffix(\"_gt\"+str(level)),\n                           (abs(pl.col(\"^responder_\\d?$\")) < level).name.suffix(\"_lt\"+str(level)+\"_abs\"),\n                           (pl.col(\"^responder_\\d?$\") < level).name.suffix(\"_lt\"+str(level)),        \n    )\n    return dataframe   \ndef aggs_daily(dataframe,by_symbol,variable_category,lag_n):\n    #Aggregating responders or features for an entire day either by symbol or for all symbols. Also has built-in daily lagging option (responders will be available for previous days).\n    if variable_category==VarCategory.RESPONDER:\n        col_expr = \"^responder_\\d?$\"\n    elif variable_category==VarCategory.FEATURE:\n        col_expr = \"^feature_\\d?\\d?$\"\n    else: \n        raise ValueError(\"Invalid var category \"+ str(variable_category))\n    if by_symbol==True:\n        suffix_tail = \"_by_sym\"+ ('_lag'+str(lag_n) if lag_n>0 else '')\n        join_cols = ['date_id','symbol_id']\n    else:\n        suffix_tail = \"_all\" + ('_lag'+str(lag_n) if lag_n>0 else '')\n        join_cols = ['date_id']        \n    df_aggs=dataframe.group_by(join_cols).agg(\n        (pl.col(col_expr)).mean().name.suffix(\"_mean\"+suffix_tail), \n        (pl.col(col_expr)).std().name.suffix(\"_std\"+suffix_tail),\n        (pl.col(col_expr)).median().name.suffix(\"_median\"+suffix_tail),\n        (pl.col(col_expr)).sum().name.suffix(\"_sum\"+suffix_tail),\n        (pl.col(col_expr)).quantile(.1).name.suffix(\"_1tile\"+suffix_tail),\n        (pl.col(col_expr)).quantile(.99).name.suffix(\"_99tile\"+suffix_tail),\n        (pl.col(col_expr)).quantile(.25).name.suffix(\"_25tile\"+suffix_tail),\n        (pl.col(col_expr)).quantile(.75).name.suffix(\"_75tile\"+suffix_tail),\n    )\n    if lag_n>0:        \n        if variable_category==VarCategory.RESPONDER:\n            col_expr = \"^responder.*$\"\n        elif variable_category==VarCategory.FEATURE:\n            col_expr = \"^feature.*$\"        \n        if by_symbol==True:\n            df_aggs=df_aggs.sort(join_cols).select(\n                pl.col('date_id'),\n                pl.col('symbol_id'),\n                (pl.col(col_expr)).shift(lag_n).over(\"symbol_id\")\n            )\n        else:           \n            df_aggs=df_aggs.sort(join_cols).select(\n                pl.col('date_id'),                \n                (pl.col(col_expr)).shift(lag_n)\n            )\n    dataframe=dataframe.join(df_aggs,on=join_cols,coalesce=True)\n    return dataframe\ndef pct_responder_lvl_daily(dataframe,by_symbol,lag_n,clean_nonlags=False):\n    #Getting daily % of each responder level. Daily lagging option built-in plus an option to drop nonlagged level flags.\n    col_expr=\"^.*_(gt|lt)\\d+(_abs)?$\"\n    if by_symbol==True:\n        suffix_tail = \"_by_sym\"+ ('_lag'+str(lag_n) if lag_n>0 else '')\n        join_cols = ['date_id','symbol_id']\n    else:\n        suffix_tail = \"_all\"+ ('_lag'+str(lag_n) if lag_n>0 else '')\n        join_cols = ['date_id']\n    df_lvl_pct=dataframe.group_by(join_cols).agg(\n        (((pl.col(col_expr)).sum())/((pl.col(col_expr)).len())).name.suffix(\"_pct\"+suffix_tail)\n    )\n    if lag_n>0:\n        col_expr=\"^.*_(gt|lt).*$\"\n        if by_symbol==True:\n            df_lvl_pct=df_lvl_pct.sort(join_cols).select(\n                pl.col('date_id'),\n                pl.col('symbol_id'),\n                (pl.col(col_expr)).shift(lag_n).over(\"symbol_id\")\n            )   \n        else:           \n            df_lvl_pct=df_lvl_pct.sort(df_lvl_pct).select(\n                pl.col('date_id'),                \n                (pl.col(col_expr)).shift(lag_n)\n            )    \n    dataframe=dataframe.join(df_lvl_pct,on=join_cols,coalesce=True)\n    if clean_nonlags==True:\n        dataframe=dataframe.drop(cs.matches(\"^.*_(gt|lt)\\d+(_abs)?$\"))\n    return dataframe\n\n#TODO: Build trailing averages\n\n\ndef feature_agg_intraday(dataframe):\n    #Averages on features intraday\n    col_expr = \"^feature_\\d?\\d?$\"\n    dataframe = dataframe.with_columns(\n        (pl.col(col_expr)).cum_sum().over(partition_by=[\"date_id\",\"symbol_id\"],order_by=[\"date_id\",\"symbol_id\"]).name.suffix(\"_cumsum_intraday\"),\n        (((pl.col(col_expr)).cum_sum().over(partition_by=[\"date_id\",\"symbol_id\"],order_by=[\"date_id\",\"symbol_id\"]))/((pl.col(col_expr)).cum_count().over(partition_by=[\"date_id\",\"symbol_id\"],order_by=[\"date_id\",\"symbol_id\"]))).name.suffix(\"_cumavg_intraday\"),\n        (pl.col(col_expr)).cum_max().over(partition_by=[\"date_id\",\"symbol_id\"],order_by=[\"date_id\",\"symbol_id\"]).name.suffix(\"_max_intraday\"),\n        (pl.col(col_expr)).cum_min().over(partition_by=[\"date_id\",\"symbol_id\"],order_by=[\"date_id\",\"symbol_id\"]).name.suffix(\"_min_intraday\"),\n     )\n    return dataframe\n    \n\ndef feature_threshold_flags(dataframe):\n    #Flag features when they go outside of some threshold based on previous day\n    calc_cols = [[x for x in dataframe.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_99tile_all_lag1' for x in dataframe.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_99tile_by_sym_lag1' for x in dataframe.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_1tile_all_lag1' for x in dataframe.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_1tile_by_sym_lag1' for x in dataframe.columns if re.match(\"^feature_\\d?\\d?$\", x)]\n                ]\n    #Can't figure out a better way to do this that actually works...\n    for idx,y in enumerate(calc_cols[0]):\n            dataframe=dataframe.with_columns(\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[1][idx])).to_series()).alias(y+\"_99tile_all_thres\"),\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[2][idx])).to_series()).alias(y+\"_99tile_by_sym_thres\"),\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[3][idx])).to_series()).alias(y+\"_1tile_all_thres\"),\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[4][idx])).to_series()).alias(y+\"_1tile_by_sym_thres\"),\n            )\n    return dataframe\n    \ndef poly2_features(dataframe,features):\n    #Create cross-multiplied features\n    for idx1,feat1 in enumerate(features):\n        for idx2,feat2 in enumerate(features):\n            if idx2>idx1:\n                dataframe=dataframe.with_columns((pl.col(feat1)*pl.col(feat2)).alias(feat1+'_'+feat2))\n    return dataframe\n\ndef apply_sine(dataframe,features):\n    #Apply sine function\n    dataframe=dataframe.with_columns((np.sin(pl.col(features))).name.suffix(\"_sin\"))\n    return dataframe","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-17T16:22:50.628121Z","iopub.execute_input":"2024-12-17T16:22:50.628523Z","iopub.status.idle":"2024-12-17T16:22:50.655732Z","shell.execute_reply.started":"2024-12-17T16:22:50.628484Z","shell.execute_reply":"2024-12-17T16:22:50.654595Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"#Test pipeline with polars lazy evaluation\npath = \"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet/\"\nq = (pl.scan_parquet(path)\n     .sort(by=['date_id','time_id'])\n#     .filter(pl.col(\"date_id\").is_in( [0,1])) #Just filtering for testing purposes\n     .filter(pl.col(\"partition_id\")==0) #Just filtering for testing purposes\n#     .sample(fraction=.1)\n \n     .pipe(responder_level_flags,1)\n     .pipe(responder_level_flags,2)\n     .pipe(responder_level_flags,3)\n     .pipe(responder_level_flags,4)\n     .pipe(responder_level_flags,5)\n\n  \n#Below calls are not working currently due to some bug in polars optimization (https://github.com/pola-rs/polars/issues/19944). \n#Getting around it for now by splitting into new pipeline.\n#    .pipe(aggs_daily,by_symbol=True,variable_category=VarCategory.RESPONDER)  \n#    .pipe(aggs_daily,by_symbol=False,variable_category=VarCategory.RESPONDER)\n#    .pipe(aggs_daily,by_symbol=True,variable_category=VarCategory.FEATURE)\n#    .pipe(aggs_daily,by_symbol=False,variable_category=VarCategory.FEATURE)    \n     \n    )\ndf = q.collect()\n\n#q2 = (df.lazy()\n#    .pipe(aggs_daily,by_symbol=True,variable_category=VarCategory.RESPONDER,lag_n=1)  \n#    .pipe(aggs_daily,by_symbol=False,variable_category=VarCategory.RESPONDER,lag_n=1)\n#    .pipe(aggs_daily,by_symbol=True,variable_category=VarCategory.FEATURE,lag_n=1)\n#    .pipe(aggs_daily,by_symbol=False,variable_category=VarCategory.FEATURE,lag_n=1)\n#    .pipe(pct_responder_lvl_daily,True,1,True)\n#    .pipe(feature_agg_intraday)\n#    )\n#df = q2.collect()\nfeat_list = df.select((pl.col(\"^feature_\\d?\\d$\"))).columns\n#df=poly2_features(df,feat_list)\ndf=apply_sine(df,feat_list)                      \n#df = feature_threshold_flags(df)\nprint(df.head())","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-17T16:22:53.442218Z","iopub.execute_input":"2024-12-17T16:22:53.442627Z","iopub.status.idle":"2024-12-17T16:22:58.205846Z","shell.execute_reply.started":"2024-12-17T16:22:53.442590Z","shell.execute_reply":"2024-12-17T16:22:58.204638Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"#print(df_samp.columns)\n\ndf_samp=df.sample(fraction=.01)\npl.Config.set_tbl_rows(2000)\n\ncorrs = []\ncorr_names = []\nfor col in [x for x in df_samp.columns if ('feature' in x) or ('responder' in x and 'lag' in x) ]:\n    corr_names.append(col)\n    corrs.append(df_samp.filter(pl.all_horizontal(pl.col(col).is_not_null())).select(pl.corr(col,\"responder_6\")).item())\n\ndf_corrs = pl.DataFrame({\"variable\": corr_names, \"corr\": corrs})\nprint(df_corrs.fill_nan(None).filter(pl.col(\"corr\").is_not_null()).sort(abs(pl.col(\"corr\")),descending=True))","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"def feature_threshold_flags(dataframe):\n    calc_cols = [[x for x in df.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_99tile_all_lag1' for x in df.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_99tile_by_sym_lag1' for x in df.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_1tile_all_lag1' for x in df.columns if re.match(\"^feature_\\d?\\d?$\", x)],\n            [x+'_1tile_by_sym_lag1' for x in df.columns if re.match(\"^feature_\\d?\\d?$\", x)]\n                ]\n    #Can't figure out a better way to do this that actually works...\n    for idx,y in enumerate(calc_cols[0]):\n            dataframe=dataframe.with_columns(\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[1][idx])).to_series()).alias(y+\"_99tile_all_thres\"),\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[2][idx])).to_series()).alias(y+\"_99tile_by_sym_thres\"),\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[3][idx])).to_series()).alias(y+\"_1tile_all_thres\"),\n                (dataframe.select(pl.col(y)).to_series() > dataframe.select(pl.col(calc_cols[4][idx])).to_series()).alias(y+\"_1tile_by_sym_thres\"),\n            )\n    return dataframe\n","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}