{ "cells": [ { "cell_type": "code", "execution_count": null, "id": "0", "metadata": {}, "outputs": [], "source": [ "%pip install \"ray[default]\"\n", "%pip install python-dotenv" ] }, { "cell_type": "code", "execution_count": null, "id": "1", "metadata": {}, "outputs": [], "source": [ "import os \n", "import ray " ] }, { "cell_type": "code", "execution_count": null, "id": "2", "metadata": {}, "outputs": [], "source": "import sys; sys.path.append(\"..\")\nimport warnings; warnings.filterwarnings(\"ignore\")\nfrom dotenv import load_dotenv; load_dotenv(override=True)\n%load_ext autoreload\n%autoreload 2" }, { "cell_type": "code", "execution_count": null, "id": "3", "metadata": {}, "outputs": [], "source": "if ray.is_initialized():\n ray.shutdown()\nray.init(\n num_cpus= 4,\n object_store_memory=2 * 1024 * 1024 * 1024,\n runtime_env={\"working_dir\": \"..\", \"env_vars\": {\"USE_LIBUV\": \"0\"}}\n)" }, { "cell_type": "code", "execution_count": null, "id": "4", "metadata": {}, "outputs": [], "source": [ "ray.cluster_resources()" ] }, { "cell_type": "code", "execution_count": null, "id": "5", "metadata": {}, "outputs": [], "source": [ "num_workers = 2\n", "resources_per_worker = {\"CPU\": 1, \"GPU\": 1}" ] }, { "cell_type": "code", "execution_count": null, "id": "6", "metadata": {}, "outputs": [], "source": [ "import os\n", "from pathlib import Path\n", "\n", "if os.path.exists(\"/efs\"):\n", " EFS_DIR = f\"/efs/shared_storage/MY-FIRST-MLOPS-PROJECT/{os.environ.get('Eva', 'default_user')}\"\n", "else: \n", " EFS_DIR = str(Path(os.getcwd()).parent / \"local_storage\")\n", " os.makedirs(EFS_DIR, exist_ok=True)\n", "\n", "print(f\"Current active storage directory: {EFS_DIR}\")" ] }, { "cell_type": "code", "execution_count": null, "id": "7", "metadata": {}, "outputs": [], "source": [ "import pandas as pd " ] }, { "cell_type": "code", "execution_count": null, "id": "8", "metadata": {}, "outputs": [], "source": "from pathlib import Path\ndataset = str(Path.cwd().parent / \"data\" / \"raw_dataset.csv\")\ndf = pd.read_csv(dataset)\ndf.head()" }, { "cell_type": "code", "execution_count": null, "id": "9", "metadata": {}, "outputs": [], "source": [ "from sklearn.model_selection import train_test_split" ] }, { "cell_type": "code", "execution_count": null, "id": "10", "metadata": {}, "outputs": [], "source": [ "df.label.value_counts()" ] }, { "cell_type": "code", "execution_count": null, "id": "11", "metadata": {}, "outputs": [], "source": [ "test_size = 0.2 \n", "train_df, val_df = train_test_split(df, stratify=df.label, test_size=test_size, random_state=42)" ] }, { "cell_type": "code", "execution_count": null, "id": "12", "metadata": {}, "outputs": [], "source": [ "train_df.label.value_counts()" ] }, { "cell_type": "code", "execution_count": null, "id": "13", "metadata": {}, "outputs": [], "source": [ "val_df.label.value_counts() * int((1 - test_size) / test_size)" ] }, { "cell_type": "code", "execution_count": null, "id": "14", "metadata": {}, "outputs": [], "source": [ "from collections import Counter \n", "import matplotlib.pyplot as plt \n", "import seaborn as sns; sns.set_theme()\n", "import warnings; warnings.filterwarnings(\"ignore\")\n", "from wordcloud import WordCloud, STOPWORDS" ] }, { "cell_type": "code", "execution_count": null, "id": "15", "metadata": {}, "outputs": [], "source": [ "tags = Counter(df.label)\n", "tags.most_common()" ] }, { "cell_type": "code", "execution_count": null, "id": "28", "metadata": {}, "outputs": [], "source": "from src.data import data, stratify_split, preprocess\nray.data.DataContext.get_current().execution_options.preserve_order = True" }, { "cell_type": "code", "execution_count": null, "id": "33", "metadata": {}, "outputs": [], "source": "import os \nimport random\nimport numpy as np\nimport torch \nfrom ray.data.preprocessor import Preprocessor\n\ndef setting_seeds(seed=42):\n np.random.seed(seed)\n random.seed(seed)\n torch.manual_seed(seed)\n torch.cuda.manual_seed(seed)\n eval(\"setattr(torch.backends.cudnn, 'deterministic', True)\") # forces GPU to always perform calculation in exact sequence\n eval(\"setattr(torch.backends.cudnn, 'benchmark', False)\") \n os.environ[\"PYTHONHASHSEED\"] = str(seed)" }, { "cell_type": "code", "execution_count": null, "id": "35", "metadata": {}, "outputs": [], "source": [ "class CustomPreprocessor():\n", " \"\"\"Custom Preprocess.\"\"\"\n", " def __init__(self, label_decoder = None):\n", " self.label_decoder = label_decoder or {\n", " 0: \"World\",\n", " 1: \"Sports\",\n", " 2: \"Business\",\n", " 3: \"Sci/Tech\"\n", " }\n", " self.class_to_index = {v:k for k, v in self.label_decoder.items()}\n", "\n", " def fitting(self, ds):\n", " _ = ds.unique(column=\"label\")\n", " return self\n", " \n", " def transforming(self, ds):\n", " return ds.map_batches(\n", " preprocess, \n", " batch_format = \"pandas\")" ] }, { "cell_type": "code", "execution_count": null, "id": "36", "metadata": {}, "outputs": [], "source": "import json\nimport torch\nimport torch.nn as nn\nfrom transformers import BertModel\nfrom transformers import AutoTokenizer" }, { "cell_type": "code", "execution_count": null, "id": "37", "metadata": {}, "outputs": [], "source": [ "model_name = \"allenai/scibert_scivocab_uncased\"\n", "tokenizer = AutoTokenizer.from_pretrained(model_name)\n", "llm = BertModel.from_pretrained(model_name, return_dict = False)\n", "embedding_dim = llm.config.hidden_size\n", "num_classes = 4" ] }, { "cell_type": "code", "execution_count": null, "id": "38", "metadata": {}, "outputs": [], "source": [ "class FinetunedLLM(nn.Module):\n", " def __init__(self, llm, dropout_p, embedding_dim, num_classes):\n", " super(FinetunedLLM, self).__init__()\n", " self.llm = llm\n", " self.dropout_p = dropout_p\n", " self.embedding_dim = embedding_dim\n", " self.num_classes = num_classes\n", " self.dropout = torch.nn.Dropout(dropout_p)\n", " self.fc1 =torch.nn.Linear(embedding_dim, num_classes)\n", "\n", " def forward(self, batch):\n", " ids, masks = batch[\"ids\"], batch[\"masks\"]\n", " seq, pool = self.llm(input_ids = ids, attention_mask = masks)\n", " z = self.dropout(pool)\n", " z = self.fc1(z)\n", " return z\n", " \n", " @torch.inference_mode()\n", " def predict(self,batch):\n", " self.eval()\n", " z = self(batch)\n", " y_pred = torch.argmax(z, dim=1).cpu().numpy()\n", " return y_pred\n", " \n", " def save(self, dp):\n", " with open(Path(dp, \"args.json\"), \"w\") as fp:\n", " contents = {\n", " \"dropout_p\": self.dropout_p,\n", " \"embedding_dim\": self.embedding_dim,\n", " \"num_classes\": self.num_classes,\n", " }\n", " json.dump(contents, fp, indent=4, sort_keys=False)\n", " torch.save(self.state_dict(), os.path.join(dp, \"model.pt\"))\n", "\n", " @classmethod\n", " def load(cls, args_fp, state_dict_fp):\n", " with open(args_fp, \"r\") as fp:\n", " kwargs = json.load(fp = fp)\n", " llm = BertModel.from_pretrained(model_name, return_dict = False)\n", " model = cls(llm = llm, **kwargs)\n", " model.load_state_dict(torch.load(state_dict_fp, map_location=torch.device(\"cpu\")))\n", " return model " ] }, { "cell_type": "code", "execution_count": null, "id": "39", "metadata": {}, "outputs": [], "source": [ "model = FinetunedLLM(llm=llm, dropout_p=0.5, embedding_dim=embedding_dim, num_classes=num_classes)\n", "print(model.named_parameters)" ] }, { "cell_type": "code", "execution_count": null, "id": "40", "metadata": {}, "outputs": [], "source": [ "from ray.train.torch import get_device" ] }, { "cell_type": "code", "execution_count": null, "id": "41", "metadata": {}, "outputs": [], "source": [ "def pad_array(arr, dtype=np.int32):\n", " max_len = max(len(row) for row in arr)\n", " padded_arr = np.zeros((arr.shape[0], max_len), dtype=dtype)\n", " for i, row in enumerate(arr):\n", " padded_arr[i][:len(row)] = row\n", " return padded_arr" ] }, { "cell_type": "code", "execution_count": null, "id": "42", "metadata": {}, "outputs": [], "source": [ "def collate_fn(batch):\n", " batch[\"ids\"] = pad_array(batch[\"ids\"])\n", " batch[\"masks\"] = pad_array(batch[\"masks\"])\n", " dtypes = {\"ids\": torch.int32, \"masks\": torch.int32, \"targets\": torch.int64}\n", " tensor_batch = {}\n", " for key, array in batch.items():\n", " tensor_batch[key] = torch.as_tensor(array, dtype=dtypes[key], device=get_device())\n", " return tensor_batch" ] }, { "cell_type": "code", "execution_count": null, "id": "43", "metadata": {}, "outputs": [], "source": [ "from pathlib import Path \n", "import ray.train as train\n", "from ray.train import Checkpoint, CheckpointConfig, DataConfig, RunConfig, ScalingConfig\n", "from ray.train.torch import TorchCheckpoint, TorchTrainer, TorchConfig\n", "import tempfile\n", "import torch.nn.functional as F\n", "from torch.nn.parallel.distributed import DistributedDataParallel" ] }, { "cell_type": "code", "execution_count": null, "id": "44", "metadata": {}, "outputs": [], "source": [ "def train_step(ds, batch_size, model, num_classes, loss_fn, optimizer):\n", " model.train()\n", " loss = 0.0\n", " ds_generator = ds.iter_torch_batches(batch_size=batch_size, collate_fn=collate_fn)\n", " for i, batch in enumerate(ds_generator):\n", " optimizer.zero_grad()\n", " z = model(batch)\n", " targets = F.one_hot(batch[\"targets\"], num_classes=num_classes).float()\n", " J = loss_fn(z, targets)\n", " J.backward()\n", " optimizer.step()\n", " loss += (J.detach().item() - loss) / (i+1)\n", " return loss\n" ] }, { "cell_type": "code", "execution_count": null, "id": "45", "metadata": {}, "outputs": [], "source": [ "def eval_step(ds, batch_size, model, num_classes, loss_fn):\n", " model.eval()\n", " loss = 0.0 \n", " y_trues, y_preds = [], []\n", " ds_generator = ds.iter_torch_batches(batch_size=batch_size, collate_fn=collate_fn)\n", " with torch.inference_mode():\n", " for i, batch in enumerate(ds_generator):\n", " z = model(batch)\n", " targets = F.one_hot(batch[\"targets\"], num_classes=num_classes).float()\n", " J = loss_fn(z, targets).item()\n", " loss += (J-loss)/(i+1)\n", " y_trues.extend(batch[\"targets\"].cpu().numpy())\n", " y_preds.extend(torch.argmax(z, dim=1).cpu().numpy())\n", " return loss, np.vstack(y_trues), np.vstack(y_preds)" ] }, { "cell_type": "code", "execution_count": null, "id": "46", "metadata": {}, "outputs": [], "source": [ "def train_loop_per_worker(config):\n", " dropout_p = config[\"dropout_p\"]\n", " lr = config[\"lr\"]\n", " lr_factor = config[\"lr_factor\"]\n", " lr_patience = config[\"lr_patience\"]\n", " num_epochs = config[\"num_epochs\"]\n", " batch_size = config[\"batch_size\"]\n", " num_classes = config[\"num_classes\"]\n", "\n", " setting_seeds()\n", " train_ds = train.get_dataset_shard(\"train\")\n", " val_ds = train.get_dataset_shard(\"val\")\n", " \n", " llm = BertModel.from_pretrained(\"allenai/scibert_scivocab_uncased\", return_dict = False)\n", " model = FinetunedLLM(llm=llm, dropout_p=dropout_p, embedding_dim=llm.config.hidden_size, num_classes=num_classes)\n", " model = train.torch.prepare_model(model)\n", "\n", " loss_fn = nn.BCEWithLogitsLoss()\n", " optimizer = torch.optim.Adam(model.parameters(), lr=lr)\n", " scheduler = torch.optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode=\"min\", factor=lr_factor, patience=lr_patience)\n", "\n", " num_workers = train.get_context().get_world_size()\n", " batch_size_per_worker = batch_size // num_workers\n", " for epoch in range(num_epochs):\n", " train_loss = train_step(train_ds, batch_size_per_worker, model, num_classes, loss_fn, optimizer)\n", " val_loss, _, _ = eval_step(val_ds, batch_size_per_worker, model, num_classes, loss_fn)\n", " scheduler.step(val_loss)\n", "\n", " with tempfile.TemporaryDirectory() as dp:\n", " if isinstance(model, DistributedDataParallel):\n", " model.module.save(dp=dp)\n", " else:\n", " model.save(dp=dp)\n", " metrics = dict(epoch=epoch, lr=optimizer.param_groups[0][\"lr\"], train_loss=train_loss, val_loss=val_loss)\n", " checkpoint = Checkpoint.from_directory(dp)\n", " train.report(metrics, checkpoint=checkpoint)" ] }, { "cell_type": "code", "execution_count": null, "id": "47", "metadata": {}, "outputs": [], "source": [ "from src.config import EFS_DIR, BASE_DIR, DATA_DIR, ARTIFACTS_DIR" ] }, { "cell_type": "code", "execution_count": null, "id": "48", "metadata": {}, "outputs": [], "source": [ "train_loop_config = {\n", " \"dropout_p\": 0.5,\n", " \"lr\": 1e-4,\n", " \"lr_factor\": 0.8,\n", " \"lr_patience\": 3,\n", " \"num_epochs\": 10,\n", " \"batch_size\": 256,\n", " \"num_classes\": num_classes,\n", "}" ] }, { "cell_type": "code", "execution_count": null, "id": "49", "metadata": {}, "outputs": [], "source": [ "scaling_config = ScalingConfig(\n", " num_workers=num_workers,\n", " use_gpu=bool(resources_per_worker[\"GPU\"]),\n", " resources_per_worker=resources_per_worker\n", ")" ] }, { "cell_type": "code", "execution_count": null, "id": "50", "metadata": {}, "outputs": [], "source": [ "checkpoint_config = CheckpointConfig(num_to_keep=1, checkpoint_score_attribute=\"val_loss\", checkpoint_score_order=\"min\")\n", "run_config = RunConfig(name=\"llm\", checkpoint_config=checkpoint_config, storage_path=EFS_DIR)" ] }, { "cell_type": "code", "execution_count": null, "id": "51", "metadata": {}, "outputs": [], "source": [ "ds = data(dataset_loc=dataset)\n", "train_ds, val_ds = stratify_split(ds, stratify=\"label\", test_size=test_size)\n", "\n", "preprocessor = CustomPreprocessor()\n", "preprocessor = preprocessor.fitting(train_ds)\n", "# materialize now, while the full cluster is free, so trainer.fit() doesn't have to\n", "# tokenize on the fly while competing with the training worker for CPU\n", "train_ds = preprocessor.transforming(train_ds).materialize()\n", "val_ds = preprocessor.transforming(val_ds).materialize()" ] }, { "cell_type": "code", "execution_count": null, "id": "52", "metadata": {}, "outputs": [], "source": [ "trainer = TorchTrainer(\n", " train_loop_per_worker=train_loop_per_worker,\n", " train_loop_config=train_loop_config,\n", " scaling_config=scaling_config,\n", " run_config=run_config,\n", " datasets={\"train\": train_ds, \"val\": val_ds},\n", " dataset_config=DataConfig(datasets_to_split=[\"train\"]),\n", " torch_config=TorchConfig(backend=\"gloo\"),\n", ")\n", "results = trainer.fit()\n", "results" ] }, { "cell_type": "code", "execution_count": null, "id": "53", "metadata": {}, "outputs": [], "source": [ "best_checkpoint = results.best_checkpoints[0][0]\n", "best_metrics = results.best_checkpoints[0][1]\n", "print(f\"best checkpoint: {best_checkpoint.path}\")\n", "print(f\"best metrics: {best_metrics}\")" ] }, { "cell_type": "code", "execution_count": null, "id": "55", "metadata": {}, "outputs": [], "source": [ "import os\n", "from huggingface_hub import create_repo, upload_folder\n", "\n", "# points directly at the epoch-0 checkpoint already saved to disk\n", "# (currently the best val_loss of any epoch trained so far)\n", "checkpoint_dir = str(EFS_DIR / \"llm\" / \"checkpoint_2026-08-10_13-03-12.020315\")\n", "\n", "HF_REPO_ID = os.environ[\"HF_REPO_ID\"] # e.g. \"yourname/agnews-scibert-classifier\", set in .env\n", "\n", "create_repo(HF_REPO_ID, exist_ok=True, repo_type=\"model\", private=True)\n", "upload_folder(\n", " folder_path=checkpoint_dir,\n", " repo_id=HF_REPO_ID,\n", " repo_type=\"model\",\n", " commit_message=\"Upload epoch 0 checkpoint (best val_loss so far)\",\n", ")\n", "print(f\"Uploaded to https://huggingface.co/{HF_REPO_ID}\")" ] }, { "cell_type": "code", "execution_count": null, "id": "57", "metadata": {}, "outputs": [], "source": [ "import os\n", "from huggingface_hub import snapshot_download\n", "\n", "HF_REPO_ID = os.environ[\"HF_REPO_ID\"]\n", "eval_checkpoint_dir = snapshot_download(repo_id=HF_REPO_ID, repo_type=\"model\")\n", "print(f\"Downloaded checkpoint to {eval_checkpoint_dir}\")" ] }, { "cell_type": "code", "execution_count": null, "id": "58", "metadata": {}, "outputs": [], "source": [ "eval_model = FinetunedLLM.load(\n", " args_fp=Path(eval_checkpoint_dir, \"args.json\"),\n", " state_dict_fp=Path(eval_checkpoint_dir, \"model.pt\"),\n", ")\n", "eval_device = torch.device(\"cuda\" if torch.cuda.is_available() else \"cpu\")\n", "eval_model = eval_model.to(eval_device)\n", "print(f\"Loaded checkpoint onto {eval_device}\")" ] }, { "cell_type": "code", "execution_count": null, "id": "59", "metadata": {}, "outputs": [], "source": "def eval_collate_fn(batch):\n batch[\"ids\"] = pad_array(batch[\"ids\"])\n batch[\"masks\"] = pad_array(batch[\"masks\"])\n dtypes = {\"ids\": torch.int32, \"masks\": torch.int32, \"targets\": torch.int64}\n return {key: torch.as_tensor(array, dtype=dtypes[key], device=eval_device) for key, array in batch.items()}\n\ndef evaluate_checkpoint(ds, batch_size, model, num_classes):\n model.eval()\n loss_fn = nn.BCEWithLogitsLoss()\n loss = 0.0\n y_trues, y_preds = [], []\n ds_generator = ds.iter_torch_batches(batch_size=batch_size, collate_fn=eval_collate_fn)\n model_device = next(model.parameters()).device\n with torch.inference_mode():\n for i, batch in enumerate(ds_generator):\n batch = {key: value.to(model_device) for key, value in batch.items()}\n z = model(batch)\n targets = F.one_hot(batch[\"targets\"], num_classes=num_classes).float()\n J = loss_fn(z, targets).item()\n loss += (J - loss) / (i + 1)\n y_trues.extend(batch[\"targets\"].cpu().numpy())\n y_preds.extend(torch.argmax(z, dim=1).cpu().numpy())\n return loss, np.array(y_trues), np.array(y_preds)" }, { "cell_type": "code", "execution_count": null, "id": "60", "metadata": {}, "outputs": [], "source": "from sklearn.metrics import classification_report\n\nval_loss, y_true, y_pred = evaluate_checkpoint(val_ds, batch_size=64, model=eval_model, num_classes=num_classes)\nprint(f\"val_loss: {val_loss:.4f}\")\nprint(classification_report(y_true, y_pred, target_names=[preprocessor.label_decoder[i] for i in sorted(preprocessor.label_decoder)]))" } ], "metadata": { "kernelspec": { "display_name": "Python 3", "language": "python", "name": "python3" }, "language_info": { "codemirror_mode": { "name": "ipython", "version": 3 }, "file_extension": ".py", "mimetype": "text/x-python", "name": "python", "nbconvert_exporter": "python", "pygments_lexer": "ipython3", "version": "3.10.11" } }, "nbformat": 4, "nbformat_minor": 5 }