{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.7.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"gpu","dataSources":[{"sourceId":38760,"databundleVersionId":4493939,"sourceType":"competition"}],"dockerImageVersionId":30302,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":true}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"**IMPORTANT NOTE: This example demonstrates how to run the synthetic data example from the [transformers4rec](https://github.com/NVIDIA-Merlin/Transformers4Rec) library. Currently, the competition data is not utilized/replaced by the synthetic data.**\n\n- https://github.com/NVIDIA-Merlin/Transformers4Rec/tree/main/examples/getting-started-session-based\n\nTransformers4Rec is a flexible and efficient library for sequential and session-based recommendation and can work with PyTorch.\n\nIt works as a bridge between NLP and recommender systems by integrating with one the most popular NLP frameworks [HuggingFace Transformers](https://github.com/huggingface/transformers), making state-of-the-art Transformer architectures available for RecSys researchers and industry practitioners.\n\n<img src=\"https://raw.githubusercontent.com/NVIDIA-Merlin/Transformers4Rec/main/_images/sequential_rec.png\" alt=\"Sequential and Session-based recommendation with Transformers4Rec\" style=\"width:800px;display:block;margin-left:auto;margin-right:auto;\"/><br>\n<div style=\"text-align: center; margin: 20pt\">\n  <figcaption style=\"font-style: italic;\">Sequential and Session-based recommendation with Transformers4Rec</figcaption>\n</div>\n\nTransformers4Rec supports multiple input features and provides configurable building blocks that can be easily combined for custom architectures.\n\nYou can build a fully GPU-accelerated pipeline for sequential and session-based recommendation with Transformers4Rec and its smooth integration with other components of [NVIDIA Merlin](https://developer.nvidia.com/nvidia-merlin):  [NVTabular](https://github.com/NVIDIA-Merlin/NVTabular) for preprocessing and [Triton Inference Server](https://github.com/triton-inference-server/server).\n\nAnd in their examples directory, you can find the following tutorials/examples:\n- [End-to-end session-based recommendation](https://github.com/NVIDIA-Merlin/Transformers4Rec/tree/main/examples/end-to-end-session-based)\n- [Tutorial - End-to-End Session-Based Recommendation on GPU](https://github.com/NVIDIA-Merlin/Transformers4Rec/tree/main/examples/tutorial)","metadata":{}},{"cell_type":"code","source":"!pip install -q transformers4rec[pytorch,nvtabular]\n!pip install -q -U nvtabular==1.3.3","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2024-10-11T11:26:10.106201Z","iopub.execute_input":"2024-10-11T11:26:10.106569Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# ETL with NVTabular","metadata":{}},{"cell_type":"markdown","source":"### Import required libraries\n","metadata":{}},{"cell_type":"code","source":"import os\nimport glob\n\nimport numpy as np\nimport pandas as pd\n\nimport cudf\nimport cupy as cp\nimport nvtabular as nvt\nfrom nvtabular.ops import *\nfrom merlin.schema.tags import Tags","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Define Input/Output Path","metadata":{}},{"cell_type":"code","source":"INPUT_DATA_DIR = os.environ.get(\"INPUT_DATA_DIR\", \"data/\")","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Create a Synthetic Input Data\n","metadata":{}},{"cell_type":"code","source":"NUM_ROWS = 1000000\nlong_tailed_item_distribution = np.clip(np.random.lognormal(3., 1., NUM_ROWS).astype(np.int32), 1, 50000)\n\n# generate random item interaction features \ndf = pd.DataFrame(np.random.randint(70000, 90000, NUM_ROWS), columns=['session_id'])\ndf['item_id'] = long_tailed_item_distribution\n\n# generate category mapping for each item-id\ndf['category'] = pd.cut(df['item_id'], bins=334, labels=np.arange(1, 335)).astype(np.int32)\ndf['timestamp/age_days'] = np.random.uniform(0, 1, NUM_ROWS)\ndf['timestamp/weekday/sin']= np.random.uniform(0, 1, NUM_ROWS)\n\n# generate day mapping for each session \nmap_day = dict(zip(df.session_id.unique(), np.random.randint(1, 10, size=(df.session_id.nunique()))))\ndf['day'] =  df.session_id.map(map_day)\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"df.head()\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Feature Engineering with NVTabular\n","metadata":{}},{"cell_type":"code","source":"# Categorify categorical features\ncateg_feats = ['session_id', 'item_id', 'category'] >> nvt.ops.Categorify(start_index=1)\n\n# Define Groupby Workflow\ngroupby_feats = categ_feats + ['day', 'timestamp/age_days', 'timestamp/weekday/sin']\n\n# Groups interaction features by session and sorted by timestamp\ngroupby_features = groupby_feats >> nvt.ops.Groupby(\n    groupby_cols=[\"session_id\"], \n    aggs={\n        \"item_id\": [\"list\", \"count\"],\n        \"category\": [\"list\"],     \n        \"day\": [\"first\"],\n        \"timestamp/age_days\": [\"list\"],\n        'timestamp/weekday/sin': [\"list\"],\n        },\n    name_sep=\"-\")\n\n# Select and truncate the sequential features\nsequence_features_truncated = (groupby_features['category-list']) >> nvt.ops.ListSlice(0,20) >> nvt.ops.Rename(postfix = '_trim')\n\nsequence_features_truncated_item = (\n    groupby_features['item_id-list']\n    >> nvt.ops.ListSlice(0,20) \n    >> nvt.ops.Rename(postfix = '_trim')\n    >> TagAsItemID()\n)  \nsequence_features_truncated_cont = (\n    groupby_features['timestamp/age_days-list', 'timestamp/weekday/sin-list'] \n    >> nvt.ops.ListSlice(0,20) \n    >> nvt.ops.Rename(postfix = '_trim')\n    >> nvt.ops.AddMetadata(tags=[Tags.CONTINUOUS])\n)\n\n# Filter out sessions with length 1 (not valid for next-item prediction training and evaluation)\nMINIMUM_SESSION_LENGTH = 2\nselected_features = (\n    groupby_features['item_id-count', 'day-first', 'session_id'] + \n    sequence_features_truncated_item +\n    sequence_features_truncated + \n    sequence_features_truncated_cont\n)\n    \nfiltered_sessions = selected_features >> nvt.ops.Filter(f=lambda df: df[\"item_id-count\"] >= MINIMUM_SESSION_LENGTH)\n\n\nworkflow = nvt.Workflow(filtered_sessions)\ndataset = nvt.Dataset(df, cpu=False)\n# Generating statistics for the features\nworkflow.fit(dataset)\n# Applying the preprocessing and returning an NVTabular dataset\nsessions_ds = workflow.transform(dataset)\n# Converting the NVTabular dataset to a Dask cuDF dataframe (`to_ddf()`) and then to cuDF dataframe (`.compute()`)\nsessions_gdf = sessions_ds.to_ddf().compute()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"sessions_gdf.head(3)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"workflow.save('workflow_etl')","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"workflow.fit_transform(dataset).to_parquet(os.path.join(INPUT_DATA_DIR, \"processed_nvt\"))","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Export pre-processed data by day\n","metadata":{}},{"cell_type":"code","source":"OUTPUT_DIR = os.environ.get(\"OUTPUT_DIR\",os.path.join(INPUT_DATA_DIR, \"sessions_by_day\"))\n!mkdir -p $OUTPUT_DIR","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"from transformers4rec.data.preprocessing import save_time_based_splits\nsave_time_based_splits(data=nvt.Dataset(sessions_gdf),\n                       output_dir= OUTPUT_DIR,\n                       partition_col='day-first',\n                       timestamp_col='session_id', \n                      )","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Checking the preprocessed outputs\n","metadata":{}},{"cell_type":"code","source":"TRAIN_PATHS = sorted(glob.glob(os.path.join(OUTPUT_DIR, \"1\", \"train.parquet\")))","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"gdf = cudf.read_parquet(TRAIN_PATHS[0])\ngdf\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# Session-based Recommendation with XLNET\n","metadata":{}},{"cell_type":"code","source":"import os\nimport rich\nimport pkg_resources\nrich.__version__ = pkg_resources.get_distribution(\"rich\").version\n\nos.environ[\"CUDA_VISIBLE_DEVICES\"]=\"0\"\n\nimport glob\nimport torch \n\nfrom transformers4rec import torch as tr\nfrom transformers4rec.torch.ranking_metric import NDCGAt, AvgPrecisionAt, RecallAt\nfrom transformers4rec.torch.utils.examples_utils import wipe_memory","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Set the schema object\n","metadata":{}},{"cell_type":"code","source":"from merlin_standard_lib import Schema\nSCHEMA_PATH = os.environ.get(\"INPUT_SCHEMA_PATH\", \"data/processed_nvt/schema.pbtxt\")\nschema = Schema().from_proto_text(SCHEMA_PATH)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"!head -20 $SCHEMA_PATH","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# You can select a subset of features for training\nschema = schema.select_by_name(['item_id-list_trim', \n                                'category-list_trim', \n                                'timestamp/weekday/sin-list_trim',\n                                'timestamp/age_days-list_trim'])\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Define the sequential input module\n","metadata":{}},{"cell_type":"code","source":"inputs = tr.TabularSequenceFeatures.from_schema(\n        schema,\n        max_sequence_length=20,\n        continuous_projection=64,\n        d_output=100,\n        masking=\"mlm\",\n)\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Define the Transformer Block\n","metadata":{}},{"cell_type":"code","source":"# Define XLNetConfig class and set default parameters for HF XLNet config  \ntransformer_config = tr.XLNetConfig.build(\n    d_model=64, n_head=4, n_layer=2, total_seq_length=20\n)\n# Define the model block including: inputs, masking, projection and transformer block.\nbody = tr.SequentialBlock(\n    inputs, tr.MLPBlock([64]), tr.TransformerBlock(transformer_config, masking=inputs.masking)\n)\n\n# Defines the evaluation top-N metrics and the cut-offs\nmetrics = [NDCGAt(top_ks=[20, 40], labels_onehot=True),  \n           RecallAt(top_ks=[20, 40], labels_onehot=True)]\n\n# Define a head related to next item prediction task \nhead = tr.Head(\n    body,\n    tr.NextItemPredictionTask(weight_tying=True, hf_format=True, \n                              metrics=metrics),\n    inputs=inputs,\n)\n\n# Get the end-to-end Model class \nmodel = tr.Model(head)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Set Training arguments\n","metadata":{}},{"cell_type":"code","source":"from transformers4rec.config.trainer import T4RecTrainingArguments\nfrom transformers4rec.torch import Trainer\n# Set hyperparameters for training \ntrain_args = T4RecTrainingArguments(data_loader_engine='nvtabular', \n                                    dataloader_drop_last = True,\n                                    gradient_accumulation_steps = 1,\n                                    per_device_train_batch_size = 128, \n                                    per_device_eval_batch_size = 32,\n                                    output_dir = \"./tmp\", \n                                    learning_rate=0.0005,\n                                    lr_scheduler_type='cosine', \n                                    learning_rate_num_cosine_cycles_by_epoch=1.5,\n                                    num_train_epochs=5,\n                                    max_sequence_length=20, \n                                    report_to = [],\n                                    logging_steps=50,\n                                    no_cuda=False)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Daily Fine-Tuning: Training over a time window\n","metadata":{}},{"cell_type":"code","source":"# Instantiate the T4Rec Trainer, which manages training and evaluation for the PyTorch API\ntrainer = Trainer(\n    model=model,\n    args=train_args,\n    schema=schema,\n    compute_metrics=True,\n)\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"INPUT_DATA_DIR = os.environ.get(\"INPUT_DATA_DIR\", \"data\")\nOUTPUT_DIR = os.environ.get(\"OUTPUT_DIR\", f\"{INPUT_DATA_DIR}/sessions_by_day\")\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"%%time\nstart_time_window_index = 1\nfinal_time_window_index = 7\n#Iterating over days of one week\nfor time_index in range(start_time_window_index, final_time_window_index):\n    # Set data \n    time_index_train = time_index\n    time_index_eval = time_index + 1\n    train_paths = glob.glob(os.path.join(OUTPUT_DIR, f\"{time_index_train}/train.parquet\"))\n    eval_paths = glob.glob(os.path.join(OUTPUT_DIR, f\"{time_index_eval}/valid.parquet\"))\n    print(train_paths)\n    \n    # Train on day related to time_index \n    print('*'*20)\n    print(\"Launch training for day %s are:\" %time_index)\n    print('*'*20 + '\\n')\n    trainer.train_dataset_or_path = train_paths\n    trainer.reset_lr_scheduler()\n    trainer.train()\n    trainer.state.global_step +=1\n    print('finished')\n    \n    # Evaluate on the following day\n    trainer.eval_dataset_or_path = eval_paths\n    train_metrics = trainer.evaluate(metric_key_prefix='eval')\n    print('*'*20)\n    print(\"Eval results for day %s are:\\t\" %time_index_eval)\n    print('\\n' + '*'*20 + '\\n')\n    for key in sorted(train_metrics.keys()):\n        print(\" %s = %s\" % (key, str(train_metrics[key]))) \n    wipe_memory()\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Saves the model\n","metadata":{}},{"cell_type":"code","source":"trainer._save_model_and_checkpoint(save_model_class=True)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Reloads the model\n","metadata":{}},{"cell_type":"code","source":"trainer.load_model_trainer_states_from_checkpoint('./tmp/checkpoint-%s'%trainer.state.global_step)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"### Re-compute eval metrics of validation data\n","metadata":{}},{"cell_type":"code","source":"eval_data_paths = glob.glob(os.path.join(OUTPUT_DIR, f\"{time_index_eval}/valid.parquet\"))","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# set new data from day 7\neval_metrics = trainer.evaluate(eval_dataset=eval_data_paths, metric_key_prefix='eval')\nfor key in sorted(eval_metrics.keys()):\n    print(\"  %s = %s\" % (key, str(eval_metrics[key])))","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"","metadata":{},"outputs":[],"execution_count":null}]}