{"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":"import pandas as pd\nimport polars as pl\nimport numpy as np\nimport torch\nimport gc\nfrom matplotlib import pyplot as plt\nimport matplotlib.cm as cm\nfrom sklearn.model_selection import StratifiedGroupKFold\nimport time\n\nimport random\nfrom tqdm import tqdm\nimport os","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T01:54:13.610083Z","iopub.execute_input":"2025-01-05T01:54:13.610361Z","iopub.status.idle":"2025-01-05T01:54:18.412293Z","shell.execute_reply.started":"2025-01-05T01:54:13.610336Z","shell.execute_reply":"2025-01-05T01:54:18.411241Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Define global variables","metadata":{}},{"cell_type":"code","source":"col_target = \"responder_6\"\ncols_feature = [f'feature_0{i}' for i in range(9)] + [f'feature_{j}' for j in range(12, 79, 1)]\ncols_responder = [f'responder_{i}' for i in range(9)]\ncols_responder_lag = [f'responder_{i}_lag_1' for i in range(9)]\ncols_category = ['feature_09', 'feature_10', 'feature_11'] # int8, int8, int16\ncols_embedding = {'feature_09': 6, 'feature_10': 6, 'feature_11': 12}\npad_feature = len(cols_feature) + sum([v for k, v in cols_embedding.items()])\npad_responder = len(cols_responder)\nstart_dt = 1100\nend_dt = 1600","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T01:55:12.973018Z","iopub.execute_input":"2025-01-05T01:55:12.973424Z","iopub.status.idle":"2025-01-05T01:55:12.981201Z","shell.execute_reply.started":"2025-01-05T01:55:12.973391Z","shell.execute_reply":"2025-01-05T01:55:12.980065Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"train_path = \"/kaggle/input/jane-street-real-time-market-data-forecasting/train.parquet\"\ntrain = pl.scan_parquet(train_path).filter((pl.col(\"date_id\")>=start_dt) & (pl.col(\"date_id\")<end_dt))\n# stats = pl.scan_parquet(train_path).filter(pl.col(\"date_id\")>=1100)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T01:55:15.137902Z","iopub.execute_input":"2025-01-05T01:55:15.138300Z","iopub.status.idle":"2025-01-05T01:55:15.170971Z","shell.execute_reply.started":"2025-01-05T01:55:15.138268Z","shell.execute_reply":"2025-01-05T01:55:15.169977Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Statistics","metadata":{}},{"cell_type":"code","source":"# def df2stats(df, cols_numerical):\n#     stats_numerical = {\n#         col: {\n#             \"min\": df[col].min(),\n#             \"min_abs\": df[col].abs().min(),\n#             \"max\": df[col].max(),\n#             \"mean\": df[col].mean(),\n#             \"std\": df[col].std()\n#         }\n#         for col in cols_numerical\n#     }\n#     return stats_numerical\n# def sanitize_df(df, cols_numerical, stats_numerical):\n#     # df.fillna(0, inplace=True)\n#     # df.replace([np.inf, -np.inf], [1, -1], inplace=True)\n#     for col in cols_numerical:\n#         pad_value = stats_numerical[col]['mean']\n#         pad_value_casted = np.array([pad_value], dtype=df[col].dtype)[0]\n#         df[col].fillna(pad_value_casted, inplace=True)\n#     # df = df.astype({col: np.float32 for col in df.select_dtypes(include=['float64']).columns})    \n#     # return df\n    \n# df = stats.collect()\n# print(df.shape)\n# stats_numerical = df2stats(df, cols_feature+cols_responder)\n# import json\n# with open('stats_numerical.json', 'w') as file:\n#     file.write(json.dumps(stats_numerical))\n# del df","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-28T09:11:59.634315Z","iopub.execute_input":"2024-12-28T09:11:59.634724Z","iopub.status.idle":"2024-12-28T09:12:31.081014Z","shell.execute_reply.started":"2024-12-28T09:11:59.634674Z","shell.execute_reply":"2024-12-28T09:12:31.080039Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# df = train.filter(pl.col(\"date_id\") == 1100).collect().to_pandas()\n# df.fillna(0, inplace=True)\n# df.head()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T01:13:48.978580Z","iopub.execute_input":"2025-01-05T01:13:48.978873Z","iopub.status.idle":"2025-01-05T01:13:50.182247Z","shell.execute_reply.started":"2025-01-05T01:13:48.978850Z","shell.execute_reply":"2025-01-05T01:13:50.181506Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Helper Functions for Trajectory Sampling","metadata":{}},{"cell_type":"code","source":"device = torch.device('cuda' if torch.cuda.is_available() else 'cpu')\nprint(f'using device: {device}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T01:55:22.181556Z","iopub.execute_input":"2025-01-05T01:55:22.181907Z","iopub.status.idle":"2025-01-05T01:55:22.191018Z","shell.execute_reply.started":"2025-01-05T01:55:22.181877Z","shell.execute_reply":"2025-01-05T01:55:22.189993Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Vectorize integers 类别映射为向量\ndef int2tensor(pos, d_model=4):\n    dim_indices = torch.arange(d_model, device=device).float()  # Shape: (d_model,)\n    angle_rates = 1 / (10000 ** (dim_indices / d_model))  # Shape: (d_model,)\n    pos_expanded = pos.unsqueeze(-1)  # Shape: pos.shape + (1,)\n    angles = pos_expanded * angle_rates  # Shape: pos.shape + (d_model,)\n    PE = torch.zeros_like(angles, device=device)\n\n    # Apply sin to even indices (0, 2, 4, ...) and cos to odd indices (1, 3, 5, ...)\n    PE[..., 0::2] = torch.sin(angles[..., 0::2])\n    PE[..., 1::2] = torch.cos(angles[..., 1::2])\n\n    return PE\n\n# Vectorize integer columns 类别列映射为向量\ndef intcol2tensor(df, cols_dim):\n    df_tensor = torch.tensor(df[list(cols_dim.keys())].values, dtype=torch.long, device=device)\n    encoded_tensors = [\n        int2tensor(df_tensor[:, i], d_model) \n        for i, d_model in enumerate(cols_dim.values())\n    ]\n    transformed_tensor = torch.cat(encoded_tensors, dim=1)\n    \n    return transformed_tensor\n\ndef df2tensor(df, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder):\n    tensor_sym = torch.tensor(df['symbol_id'].values, dtype=int, device=device)\n    tensor_tim = torch.tensor(df['time_id'].values, dtype=int, device=device)\n    uniq_tim = torch.unique(tensor_tim)\n    pad_time = len(uniq_tim)\n    \n    # Create the padded 2D vector\n    tensor_fea = torch.zeros((pad_time, pad_symbols, pad_feature), device=device)\n    tensor_rsp = torch.zeros((pad_time, pad_symbols, pad_responder), device=device)\n    tensor_wgh = torch.zeros((pad_time, pad_symbols), device=device)\n    tensor_fea[tensor_tim, tensor_sym, :len(cols_feature)] = torch.tensor(df[cols_feature].values, device=device)\n    tensor_fea[tensor_tim, tensor_sym, len(cols_feature):] = intcol2tensor(df, cols_category)\n\n    tensor_rsp[tensor_tim, tensor_sym, :] = torch.tensor(df[cols_responder].values, device=device)\n\n    tensor_wgh[tensor_tim, tensor_sym] = torch.tensor(df['weight'].values, device=device)\n    return tensor_fea, tensor_rsp, tensor_wgh","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T02:10:33.471367Z","iopub.execute_input":"2025-01-05T02:10:33.471704Z","iopub.status.idle":"2025-01-05T02:10:33.482539Z","shell.execute_reply.started":"2025-01-05T02:10:33.471674Z","shell.execute_reply":"2025-01-05T02:10:33.481440Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"f_list = []\nr_list = []\nw_list = []\nfor i in range(start_dt, end_dt):\n    print(i)\n    df = train.filter(pl.col(\"date_id\") == i).collect().to_pandas()\n    df.fillna(0, inplace=True)\n    f,r,w = df2tensor(df, cols_feature, cols_embedding, cols_responder, 42, pad_feature, pad_responder)\n    f_list.append(f.to(torch.float32))\n    r_list.append(r.to(torch.float32))\n    w_list.append(w.to(torch.float32))","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T02:19:39.330395Z","iopub.execute_input":"2025-01-05T02:19:39.331229Z","iopub.status.idle":"2025-01-05T02:19:41.980377Z","shell.execute_reply.started":"2025-01-05T02:19:39.331173Z","shell.execute_reply":"2025-01-05T02:19:41.979434Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"f_tensor = torch.stack(f_list)\nr_tensor = torch.stack(r_list)\nw_tensor = torch.stack(w_list)\n\ntorch.save(f_tensor, f'dataset_fea.pt')\ntorch.save(r_tensor, f'dataset_rsp.pt')\ntorch.save(w_tensor, f'dataset_wgh.pt')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2025-01-05T02:19:44.964568Z","iopub.execute_input":"2025-01-05T02:19:44.964922Z","iopub.status.idle":"2025-01-05T02:19:45.208727Z","shell.execute_reply.started":"2025-01-05T02:19:44.964892Z","shell.execute_reply":"2025-01-05T02:19:45.207740Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"\n\n# # intcol2vec(df.iloc[:5], cols_category).shape\n\n# # Extract symbol vector from time-group into dense format\n# # 从time-group的一组symbol列提取致密2D向量 (num rows x num symbol columns)\n# def symbols2vec(df, cols_feature, cols_category, cols_responder):\n#     # df is a time_id group\n#     # Initialize a 2D numpy array for the result\n#     vec_feature = None\n#     vec_responder = None\n\n#     # Process feature columns\n#     if not cols_feature is None:\n#         vec_feature = df[cols_feature].to_numpy()\n\n#         vec_category = intcol2vec(df, cols_category)\n#         vec_feature = np.hstack((vec_feature, vec_category))\n\n#     # Process responder columns\n#     if not cols_responder is None:\n#         vec_responder = df[cols_responder].to_numpy()\n\n#     return vec_feature, vec_responder\n\n# # Put dense symbol vectors into sparse vector\n# # time-group映射为1D稀疏向量 (填充symbols并压缩该维度)\n# def time2vec(df, time_id, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder):\n#     # df is a time_id group\n#     dense_feature, dense_responder = symbols2vec(df, cols_feature, cols_category, cols_responder)\n        \n#     sparse_feature = np.zeros((pad_symbols, pad_feature))\n#     sparse_responder = np.zeros((pad_symbols, pad_responder))\n#     sparse_weight = np.zeros(pad_symbols)\n    \n#     for i, row in enumerate(df.itertuples(index=False)):\n#         symbol_id = int(row.symbol_id)\n#         if symbol_id >= pad_symbols:\n#             break\n#         if not dense_feature is None:\n#             sparse_feature[symbol_id] = dense_feature[i]\n#         if not dense_responder is None:\n#             sparse_responder[symbol_id] = dense_responder[i]\n#         sparse_weight[symbol_id] = row.weight\n\n#     return sparse_feature.flatten(), sparse_responder.flatten(), sparse_weight\n\n# # time2vec(df.iloc[:5], 0, cols_feature, cols_category, None, 64, 100).shape\n\n# # Sample a trajectory from a dataframe. Sample 'root length' number of points using chebyshev nodes\n# def df2traj_time(df, traj_len, eplison, include_last=False, rng=None):\n#     # df is a date_id group\n#     if rng is None:\n#         rng = np.random.default_rng(seed=42)\n#     time_uniq = list(df['time_id'].unique())\n\n#     chebyshev_nodes = np.array([np.cos((2 * i + 1) / (2 * traj_len) * np.pi + rng.uniform(-eplison, eplison)) for i in range(traj_len)])\n#     chebyshev_nodes = (chebyshev_nodes + 1) / 2 * len(time_uniq)\n#     chebyshev_nodes.sort()\n#     time_traj = [time_uniq[int(i)] for i in chebyshev_nodes]\n#     if eplison < 0.01 or include_last:\n#         time_traj[-1] = time_uniq[-1]\n#     return time_traj\n# # df2traj(df[:5], 5)\n\n\n# def df2traj_vec(df, time_traj, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time):\n#     # df is a date_id group\n#     groups = df.groupby('time_id')\n    \n#     # Create the padded 2D vector\n#     traj_fea = np.zeros((pad_time, pad_symbols * pad_feature))\n#     traj_rsp = np.zeros((pad_time, pad_symbols * pad_responder))\n#     traj_wgh = np.zeros((pad_time, pad_symbols))\n\n#     for i, time_id in enumerate(time_traj):\n#         if i >= pad_time:\n#             break\n#         time_group = groups.get_group(time_id)\n#         vec_fea, vec_rsp, vec_wgh = time2vec(time_group, time_id, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder)\n        \n#         # time_embedding = int2vec(time_id+time_shift, pad_symbols * (pad_feature + pad_responder)).flatten()\n#         # time_vec += time_embedding\n#         traj_fea[i] = vec_fea\n#         traj_rsp[i] = vec_rsp\n#         traj_wgh[i] = vec_wgh\n        \n#     return traj_fea, traj_rsp, traj_wgh","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-28T09:12:31.088516Z","iopub.execute_input":"2024-12-28T09:12:31.088900Z","iopub.status.idle":"2024-12-28T09:12:31.109392Z","shell.execute_reply.started":"2024-12-28T09:12:31.088874Z","shell.execute_reply":"2024-12-28T09:12:31.108245Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# import torch\n# from multiprocessing import Pool\n# import time\n# import polars as pl\n# from tqdm import tqdm\n# import psutil\n\n# def process_worker(args):\n#     df_lag, df_curr, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time, num_runs, rng = args\n#     if rng is None:\n#         rng = np.random.default_rng(seed=42)\n#     list_fea = []\n#     list_rsp = []\n#     list_wgh = []\n#     list_tim = []\n    \n#     time_shift = -1 * df_lag['time_id'].max()\n#     for _ in range(num_runs):\n#         traj_lag = df2traj_time(df_lag, pad_time // 2, 0.5, False, rng) # sample with random\n#         vec_lag_fea, vec_lag_rsp, vec_lag_wgh = df2traj_vec(df_lag, traj_lag, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time//2)\n        \n#         time_now = rng.integers(0, df_curr['time_id'].max())\n#         df_curr_hist = df_curr.loc[df_curr['time_id']<=time_now]\n#         traj_curr = df2traj_time(df_curr_hist, pad_time // 2, 0.1, True, rng) # sample without random, must include last\n#         vec_cur_fea, vec_cur_rsp, vec_cur_wgh = df2traj_vec(df_curr_hist, traj_curr, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time//2)\n#         vec_fea = np.vstack((vec_lag_fea, vec_cur_fea))\n#         vec_rsp = np.vstack((vec_lag_rsp, vec_cur_rsp))\n#         vec_wgh = np.vstack((vec_lag_wgh, vec_cur_wgh))\n#         traj_lag = time_shift*np.ones(pad_time//2) + traj_lag\n#         vec_tim = np.hstack((traj_lag, traj_curr))\n\n#         list_fea.append(torch.tensor(vec_fea, dtype=torch.float16))\n#         list_rsp.append(torch.tensor(vec_rsp, dtype=torch.float16))\n#         list_wgh.append(torch.tensor(vec_wgh, dtype=torch.float16))\n#         list_tim.append(torch.tensor(vec_tim, dtype=torch.int16))\n\n#     return list_fea, list_rsp, list_wgh, list_tim\n    \n# def df2dataset(ldf, stats_numerical, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time, runs, num_workers=-1, verbose=False):\n#     # Get unique `date_id` values and split them\n#     date_split = ldf.select(\"date_id\").unique().collect()['date_id'].to_list()\n#     date_split.sort()\n#     list_fea = []\n#     list_rsp = []\n#     list_wgh = []\n#     list_tim = []\n\n#     with Pool(processes=4) as pool:\n#         for i in tqdm(range(len(date_split)-1)):\n#             if verbose:\n#                 print(f'{i} / {len(date_split)}')\n#                 memory = psutil.virtual_memory()\n#                 print(f\"Total: {memory.total / 1e9:.2f} GB\")\n#                 print(f\"Available: {memory.available / 1e9:.2f} GB\")\n#                 print(f\"Used: {memory.used / 1e9:.2f} GB\")\n#                 print(f\"Percentage: {memory.percent}%\")\n#             date_lag = date_split[i]\n#             date_curr = date_split[i+1]\n#             df_lag = ldf.filter(pl.col(\"date_id\") == date_lag).collect().to_pandas()\n#             df_curr = ldf.filter(pl.col(\"date_id\") == date_curr).collect().to_pandas()\n#             sanitize_df(df_lag, cols_feature+cols_responder, stats_numerical)\n#             sanitize_df(df_curr, cols_feature+cols_responder, stats_numerical)\n            \n#             # df is a date_id group\n#             arguements = [df_lag, df_curr, cols_feature, cols_category, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time, runs]\n#             if num_workers > 1:\n#             # use parallel\n#                 tasks = [arguements + [np.random.default_rng(seed=worker_idx+(int(time.time() * 1e6) % (2**32)))] for worker_idx in range(num_workers)]\n#                 # Execute tasks in parallel\n#                 results = pool.map(process_worker, tasks)\n#                 # Combine results from all workers\n#                 for result in results:\n#                     list_fea.extend(result[0])\n#                     list_rsp.extend(result[1])\n#                     list_wgh.extend(result[2])\n#                     list_tim.extend(result[3])\n#             else:\n#             # use sequential\n#                 part_fea, part_rsp, part_wgh, part_tim = process_worker(arguements+[None])\n#                 list_fea.extend(part_fea)\n#                 list_rsp.extend(part_rsp)\n#                 list_wgh.extend(part_wgh)\n#                 list_tim.extend(part_tim)\n#             if verbose:\n#                 print(f'check: {list_tim[-1].sum() % 960}')\n#             if i > 0 and (i+1) % 50 == 0:\n#                 split = i // 50\n#                 print(f'completed split{split}, size {len(list_fea)}')\n                \n#                 tensor_fea = torch.stack(list_fea)\n#                 list_fea = []\n#                 torch.save(tensor_fea, f'dataset_fea_{i}.pt')\n#                 del tensor_fea\n\n#                 tensor_rsp = torch.stack(list_rsp)\n#                 list_rsp = []\n#                 torch.save(tensor_rsp, f'dataset_rsp_{i}.pt')\n#                 del tensor_rsp\n                \n#                 tensor_wgh = torch.stack(list_wgh)\n#                 list_wgh = []\n#                 torch.save(tensor_wgh, f'dataset_wgh_{i}.pt')\n#                 del tensor_wgh\n\n#                 tensor_tim = torch.stack(list_tim)\n#                 list_tim = []\n#                 torch.save(tensor_tim, f'dataset_tim_{i}.pt')\n#                 del tensor_tim","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-28T09:12:31.110483Z","iopub.execute_input":"2024-12-28T09:12:31.110833Z","iopub.status.idle":"2024-12-28T09:12:31.131912Z","shell.execute_reply.started":"2024-12-28T09:12:31.110798Z","shell.execute_reply":"2024-12-28T09:12:31.130979Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# pad_symbols = 42\n# pad_time = 16\n# runs = 32\n# df2dataset(train, stats_numerical, cols_feature, cols_embedding, cols_responder, pad_symbols, pad_feature, pad_responder, pad_time, runs, num_workers=4, verbose=True)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-12-28T09:12:31.132811Z","iopub.execute_input":"2024-12-28T09:12:31.133238Z","execution_failed":"2024-12-28T09:13:56.804Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# fea = torch.load('/kaggle/working/dataset_fea_1.pt')\n# rsp = torch.load('/kaggle/working/dataset_rsp_1.pt')\n# tim = torch.load('/kaggle/working/dataset_tim_1.pt')\n# wgh = torch.load('/kaggle/working/dataset_wgh_1.pt')\n# tim[2, 15]\n# fea[2, 15, pad_feature:2*pad_feature]\n# rsp[2, 15, pad_responder:2*pad_responder]\n# wgh[2, 15]\n# df = train.filter(pl.col(\"date_id\") == 1102).collect().to_pandas()\n# print(df['time_id'].max())\n# df.loc[df['time_id']==735]","metadata":{"trusted":true,"execution":{"execution_failed":"2024-12-28T09:13:56.805Z"}},"outputs":[],"execution_count":null}]}