{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.11.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":99552,"databundleVersionId":13851420,"sourceType":"competition"}],"dockerImageVersionId":31089,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"import os\nimport shutil\nimport numpy as np\nimport pandas as pd\nimport cv2\nimport pydicom\nimport torch\nimport torch.nn as nn\nimport timm\nimport albumentations as A\nfrom albumentations.pytorch import ToTensorV2\nfrom torch.utils.data import Dataset, DataLoader\nfrom sklearn.model_selection import StratifiedKFold\nfrom sklearn.metrics import roc_auc_score\nfrom tqdm import tqdm\nfrom joblib import Parallel, delayed # Key imports for parallel processing\nimport polars as pl\nimport kaggle_evaluation.rsna_inference_server\nfrom IPython.display import display\n\n# ----------------------- DEVICE & CONSTANTS SETUP ------------------------\ndevice = torch.device('cuda' if torch.cuda.is_available() else 'cpu')\nprint(f\"Using device: {device}\")\n\nif device.type == 'cuda':\n    torch.backends.cudnn.benchmark = True\n    scaler = torch.cuda.amp.GradScaler()\n\n# Global Constants\nIMAGE_SIZE = 256\nNUM_SLICES = 32\nBATCH_SIZE = 32\nNUM_FOLDS = 3\nEPOCHS = 30\nNUM_WORKERS = 8      # Workers for PyTorch DataLoader (I/O acceleration during training)\nNUM_CPU_CORES = 6    # Cores for Joblib (Preprocessing acceleration)\n\nEARLY_STOPPING_PATIENCE = 5\nMIN_DELTA = 1e-4\n\nID_COL = 'SeriesInstanceUID'\nLABEL_COLS = [\n    'Left Infraclinoid Internal Carotid Artery','Right Infraclinoid Internal Carotid Artery',\n    'Left Supraclinoid Internal Carotid Artery','Right Supraclinoid Internal Carotid Artery',\n    'Left Middle Cerebral Artery','Right Middle Cerebral Artery','Anterior Communicating Artery',\n    'Left Anterior Cerebral Artery','Right Anterior Cerebral Artery','Left Posterior Communicating Artery',\n    'Right Posterior Communicating Artery','Basilar Tip','Other Posterior Circulation','Aneurysm Present'\n]\nANEURYSM_IDX = LABEL_COLS.index('Aneurysm Present')\n\nTRAIN_CSV = \"/kaggle/input/rsna-intracranial-aneurysm-detection/train.csv\"\nTRAIN_LOCALIZERS = \"/kaggle/input/rsna-intracranial-aneurysm-detection/train_localizers.csv\"\nSERIES_PATH = \"/kaggle/input/rsna-intracranial-aneurysm-detection/series\"\n\n# Global Model Ensemble and Transforms (Will be lazily loaded for inference)\nENSEMBLE_MODELS = None\nVAL_TRANSFORM = None \n\n# ----------------------- AUGMENTATIONS ------------------------\ntrain_transform = A.Compose([\n    A.Resize(IMAGE_SIZE, IMAGE_SIZE),\n    A.HorizontalFlip(p=0.5),\n    A.VerticalFlip(p=0.5),\n    A.RandomRotate90(p=0.5),\n    A.ShiftScaleRotate(shift_limit=0.0625, scale_limit=0.1, rotate_limit=15, p=0.7),\n    A.ElasticTransform(alpha=1, sigma=50, alpha_affine=50, p=0.3),\n    A.RandomBrightnessContrast(brightness_limit=0.2, contrast_limit=0.2, p=0.5),\n    A.Normalize(mean=[0.5], std=[0.5]),\n    ToTensorV2()\n])\n\nval_transform_definition = A.Compose([\n    A.Resize(IMAGE_SIZE, IMAGE_SIZE),\n    A.Normalize(mean=[0.5], std=[0.5]),\n    ToTensorV2()\n])\n\n# ----------------------- DICOM PREPROCESSING FUNCTION ------------------------\ndef process_dicom_series(series_path):\n    \"\"\"Reads DICOM series and returns a single-channel grayscale 2D MIP.\"\"\"\n    try:\n        files = [os.path.join(series_path, f) for f in os.listdir(series_path) if f.lower().endswith(\".dcm\")]\n        if len(files) == 0:\n            raise ValueError(f\"No DICOM files in {series_path}\")\n        slices = []\n        for f in files:\n            try:\n                ds = pydicom.dcmread(f, force=True)\n                z = getattr(ds, 'ImagePositionPatient', [0,0,0])[2]\n                arr = ds.pixel_array\n                slices.append((z, arr))\n            except Exception:\n                continue\n        if len(slices) == 0:\n            raise ValueError(f\"No valid slices in {series_path}\")\n        slices.sort(key=lambda x: x[0])\n        volume = np.stack([s[1] for s in slices], axis=0)\n        \n        # Intensity windowing & normalization\n        vmin, vmax = np.percentile(volume, 5), np.percentile(volume, 95)\n        if vmax == vmin:\n            volume = (volume / (volume.max() or 1) * 255.0).astype(np.uint8)\n        else:\n            volume = np.clip(volume, vmin, vmax)\n            volume = ((volume - volume.min()) / (volume.max() - volume.min()) * 255.0).astype(np.uint8)\n            \n        resized_slices = [cv2.resize(s, (IMAGE_SIZE, IMAGE_SIZE), interpolation=cv2.INTER_LINEAR) for s in volume]\n        volume_resized = np.stack(resized_slices, axis=0)\n        \n        # Sample NUM_SLICES\n        if len(volume_resized) >= NUM_SLICES:\n            idxs = np.linspace(0, len(volume_resized)-1, NUM_SLICES).astype(int)\n            volume_resized = volume_resized[idxs]\n        else:\n            pad = NUM_SLICES - len(volume_resized)\n            volume_resized = np.pad(volume_resized, ((0,pad),(0,0),(0,0)), mode='edge')\n            \n        # Compute MIP (Max Intensity Projection)\n        start, end = int(NUM_SLICES*0.15), max(int(NUM_SLICES*0.85), 1)\n        mip = np.max(volume_resized[start:end], axis=0)\n        \n        # Final normalization to 0-255\n        range_val = mip.max() - mip.min()\n        img_out = ((mip - mip.min()) / (range_val or 1) * 255.0).astype(np.uint8)\n        return img_out[None, :, :]\n    except Exception:\n        return np.zeros((1, IMAGE_SIZE, IMAGE_SIZE), dtype=np.uint8)\n\n\n# ----------------------- HELPER FUNCTION FOR JOBILB ------------------------\nPREPROCESSED_TRAIN_DIR = \"/kaggle/working/preprocessed_train/\"\n\ndef preprocess_and_save_series(series_id, series_base_path, output_dir):\n    \"\"\"Worker function for Joblib to preprocess and save a single series.\"\"\"\n    out_path = os.path.join(output_dir, f\"{series_id}.npy\")\n    if not os.path.exists(out_path):\n        img = process_dicom_series(os.path.join(series_base_path, series_id))\n        np.save(out_path, img)\n    return True\n\n# ----------------------- PRE-CALCULATE DICOM SERIES (PARALLELIZED) ------------------------\ndf = pd.read_csv(TRAIN_CSV)\nlocalizers = pd.read_csv(TRAIN_LOCALIZERS)\ndf = df[df[ID_COL].isin(localizers[ID_COL].unique())].reset_index(drop=True)\n\nos.makedirs(PREPROCESSED_TRAIN_DIR, exist_ok=True)\n\nprint(f\"Starting PARALLEL DICOM preprocessing using {NUM_CPU_CORES} cores...\")\n\nseries_list = df[ID_COL].unique()\n\n# Parallel execution using joblib: dramatically speeds up this I/O and CPU-bound task.\nParallel(n_jobs=NUM_CPU_CORES, verbose=10)(\n    delayed(preprocess_and_save_series)(\n        series_id, \n        SERIES_PATH, \n        PREPROCESSED_TRAIN_DIR\n    ) \n    for series_id in series_list\n)\n\nprint(\"Preprocessing complete (Parallelized).\")\n\n# ----------------------- DATASET CLASSES ------------------------\nclass RSNADataset(Dataset):\n    def __init__(self, df, preprocessed_dir, transform=None):\n        self.df = df\n        self.dir = preprocessed_dir\n        self.transform = transform\n        self.labels = df[LABEL_COLS].values.astype(np.float32)\n\n    def __len__(self): return len(self.df)\n\n    def __getitem__(self, idx):\n        row = self.df.iloc[idx]\n        img = np.load(os.path.join(self.dir, f\"{row[ID_COL]}.npy\"))\n        \n        if self.transform:\n            img = self.transform(image=img.transpose(1,2,0))['image']\n            \n        labels = torch.tensor(self.labels[idx])\n        return img, labels\n\n# ----------------------- MODEL DEFINITION ------------------------\nclass EfficientNetMultiLabel(nn.Module):\n    def __init__(self, num_classes=len(LABEL_COLS)):\n        super().__init__()\n        # Load EfficientNetV2-S backbone, specifying single channel (in_chans=1)\n        self.backbone = timm.create_model(\n            'tf_efficientnetv2_s', \n            pretrained=False, \n            in_chans=1, \n            num_classes=0\n        )\n        self.pool = nn.AdaptiveAvgPool2d(1)\n        num_features = self.backbone.num_features\n        \n        # Multi-layer classifier (Head)\n        self.classifier = nn.Sequential(\n            nn.Linear(num_features, 512),\n            nn.ReLU(),\n            nn.Dropout(0.3),\n            nn.Linear(512, 256),\n            nn.ReLU(),\n            nn.Dropout(0.3),\n            nn.Linear(256, num_classes)\n        )\n\n    def forward(self, x):\n        x = self.backbone(x)\n        if x.ndim == 4:\n            x = self.pool(x).flatten(1)\n        elif x.ndim == 3:\n            x = x.mean(dim=1)\n        return self.classifier(x)\n\n# ----------------------- CLASS WEIGHTS & LOSS FUNCTIONS ------------------------\n# Calculate class weights for BCE loss\npositive_counts = df[LABEL_COLS].sum(axis=0)\ntotal_samples = len(df)\nnegative_counts = total_samples - positive_counts\nclass_weights = negative_counts / positive_counts.clip(lower=1)\nPOS_WEIGHT_TENSOR = torch.tensor(class_weights.values.astype(np.float32), dtype=torch.float32).to(device)\n\ndef dice_loss(preds, targets, smooth=1e-6):\n    \"\"\"Calculates soft Dice Loss.\"\"\"\n    preds = torch.sigmoid(preds)\n    intersection = (preds * targets).sum(dim=1)\n    return 1 - ((2 * intersection + smooth) / (preds.sum(dim=1) + targets.sum(dim=1) + smooth)).mean()\n\ndef combined_loss(preds, targets):\n    \"\"\"Combined BCE with weighted positive samples and Dice loss.\"\"\"\n    bce = nn.BCEWithLogitsLoss(pos_weight=POS_WEIGHT_TENSOR)\n    return bce(preds, targets) + dice_loss(preds, targets)\n\n# ----------------------- TRAINING LOOP ------------------------\nskf = StratifiedKFold(n_splits=NUM_FOLDS, shuffle=True, random_state=42)\nfolds = list(skf.split(df, df['Aneurysm Present']))\n\nfor fold_idx, (train_idx, val_idx) in enumerate(folds):\n    print(f\"\\n========== Fold {fold_idx} ==========\")\n    train_df, val_df = df.iloc[train_idx], df.iloc[val_idx]\n\n    train_dataset = RSNADataset(train_df, PREPROCESSED_TRAIN_DIR, transform=train_transform)\n    val_dataset = RSNADataset(val_df, PREPROCESSED_TRAIN_DIR, transform=val_transform_definition)\n\n    train_loader = DataLoader(train_dataset, batch_size=BATCH_SIZE, shuffle=True, num_workers=NUM_WORKERS, pin_memory=True)\n    val_loader = DataLoader(val_dataset, batch_size=BATCH_SIZE, shuffle=False, num_workers=NUM_WORKERS, pin_memory=True)\n\n    model = EfficientNetMultiLabel().to(device)\n    optimizer = torch.optim.AdamW(model.parameters(), lr=1e-3)\n    scheduler = torch.optim.lr_scheduler.ReduceLROnPlateau(optimizer, mode='max', factor=0.5, patience=3)\n\n    best_score = -1\n    epochs_without_improvement = 0\n\n    for epoch in range(1, EPOCHS+1):\n        # -------- TRAIN --------\n        model.train()\n        train_loss = 0\n        train_preds, train_labels = [], []\n\n        for imgs, labels in tqdm(train_loader, desc=f\"Fold {fold_idx} Epoch {epoch} Train\"):\n            imgs, labels = imgs.to(device), labels.to(device)\n            optimizer.zero_grad()\n\n            if device.type == 'cuda':\n                # Automatic Mixed Precision (AMP)\n                with torch.cuda.amp.autocast():\n                    outputs = model(imgs)\n                    loss = combined_loss(outputs, labels)\n                scaler.scale(loss).backward()\n                scaler.step(optimizer)\n                scaler.update()\n            else:\n                outputs = model(imgs)\n                loss = combined_loss(outputs, labels)\n                loss.backward()\n                optimizer.step()\n\n            train_loss += loss.item() * imgs.size(0)\n            train_preds.append(torch.sigmoid(outputs).detach().cpu())\n            train_labels.append(labels.detach().cpu())\n\n        train_loss /= len(train_dataset)\n        train_preds = torch.cat(train_preds)\n        train_labels = torch.cat(train_labels)\n        \n        # --- AUC-ROC Metric Calculation ---\n        aucs = []\n        for i in range(len(LABEL_COLS)):\n            labels_i = train_labels[:, i].numpy()\n            if len(np.unique(labels_i)) == 2:\n                aucs.append(roc_auc_score(labels_i, train_preds[:, i].numpy()))\n            else:\n                aucs.append(0.5)\n        \n        train_score = 0.5 * (aucs[ANEURYSM_IDX] + np.mean([aucs[i] for i in range(len(LABEL_COLS)) if i != ANEURYSM_IDX]))\n        print(f\"Epoch {epoch} Train Loss {train_loss:.4f} Score {train_score:.4f}\")\n\n        # -------- VALIDATION --------\n        model.eval()\n        val_loss = 0\n        val_preds, val_labels = [], []\n        with torch.no_grad():\n            for imgs, labels in tqdm(val_loader, desc=f\"Fold {fold_idx} Epoch {epoch} Val\"):\n                imgs, labels = imgs.to(device), labels.to(device)\n                \n                if device.type == 'cuda':\n                    with torch.cuda.amp.autocast():\n                        outputs = model(imgs)\n                        loss = combined_loss(outputs, labels)\n                else:\n                    outputs = model(imgs)\n                    loss = combined_loss(outputs, labels)\n\n                val_loss += loss.item() * imgs.size(0)\n                val_preds.append(torch.sigmoid(outputs).cpu())\n                val_labels.append(labels.cpu())\n\n        val_loss /= len(val_dataset)\n        val_preds = torch.cat(val_preds)\n        val_labels = torch.cat(val_labels)\n        \n        # --- AUC-ROC Metric Calculation for Validation ---\n        val_aucs = []\n        for i in range(len(LABEL_COLS)):\n            labels_i = val_labels[:, i].numpy()\n            if len(np.unique(labels_i)) == 2:\n                val_aucs.append(roc_auc_score(labels_i, val_preds[:, i].numpy()))\n            else:\n                val_aucs.append(0.5)\n\n        val_score = 0.5 * (val_aucs[ANEURYSM_IDX] + np.mean([val_aucs[i] for i in range(len(LABEL_COLS)) if i != ANEURYSM_IDX]))\n        print(f\"Epoch {epoch} Val Loss {val_loss:.4f} Score {val_score:.4f}\")\n\n        scheduler.step(val_score)\n\n        # Early stopping and checkpoint saving\n        if val_score > best_score + MIN_DELTA:\n            best_score = val_score\n            epochs_without_improvement = 0\n            torch.save(model.state_dict(), f\"best_model_fold{fold_idx}.pth\")\n            print(f\"✅ Checkpoint saved at Epoch {epoch} with score: {best_score:.4f}\")\n        else:\n            epochs_without_improvement += 1\n            print(f\"Patience: {epochs_without_improvement}/{EARLY_STOPPING_PATIENCE}\")\n            if epochs_without_improvement >= EARLY_STOPPING_PATIENCE:\n                print(f\"🛑 Early stopping triggered at Epoch {epoch}.\")\n                break\n\n\n# ----------------------- INFERENCE FUNCTION SETUP ------------------------\n\ndef load_ensemble_models():\n    \"\"\"Loads trained models into a global list (Lazy Loading).\"\"\"\n    models = []\n    global VAL_TRANSFORM\n    VAL_TRANSFORM = val_transform_definition\n    try:\n        for fold_idx in range(NUM_FOLDS):\n            model = EfficientNetMultiLabel().to(device)\n            model.load_state_dict(torch.load(f\"best_model_fold{fold_idx}.pth\", map_location=device))\n            model.eval()\n            models.append(model)\n        print(f\"Loaded {len(models)} models for ensembling.\")\n    except FileNotFoundError:\n        print(\"WARNING: Could not load trained models. Using randomly initialized models for the API.\")\n        for _ in range(NUM_FOLDS):\n            model = EfficientNetMultiLabel().to(device)\n            model.eval()\n            models.append(model)\n    return models\n\ndef predict(series_path: str) -> pl.DataFrame:\n    \"\"\"\n    Kaggle API inference function. Predicts labels for a single DICOM series path.\n    Must return a Polars DataFrame containing ONLY the prediction columns.\n    \"\"\"\n    global ENSEMBLE_MODELS, VAL_TRANSFORM\n\n    # --- 1. Lazy Model Loading (runs only on the first call) ---\n    if ENSEMBLE_MODELS is None:\n        ENSEMBLE_MODELS = load_ensemble_models()\n\n    if not ENSEMBLE_MODELS:\n        # Emergency Fallback\n        data = {col: [0.5] for col in LABEL_COLS}\n        return pl.DataFrame(data)\n\n    # --- 2. Preprocessing & Transformation ---\n    img_np = process_dicom_series(series_path) \n    img_tensor = VAL_TRANSFORM(image=img_np.transpose(1, 2, 0))['image'].to(device)\n    img_tensor = img_tensor.unsqueeze(0) # Add batch dimension (N=1)\n\n    # --- 3. Ensemble Prediction ---\n    ensemble_probs = []\n    with torch.no_grad():\n        for model in ENSEMBLE_MODELS:\n            if device.type == 'cuda':\n                with torch.cuda.amp.autocast():\n                    output = model(img_tensor)\n            else:\n                output = model(img_tensor)\n                \n            prob = torch.sigmoid(output).cpu().numpy().flatten()\n            ensemble_probs.append(prob)\n\n    # Mean of all fold predictions\n    final_prediction = np.mean(ensemble_probs, axis=0).astype(np.float64)\n    \n    # --- 4. Conversion to Polars DataFrame (REQUIRED FORMAT) ---\n    data = {col: [val] for col, val in zip(LABEL_COLS, final_prediction)}\n    \n    # CRITICAL: Clean up shared directory to prevent \"out of disk space\" errors\n    shutil.rmtree('/kaggle/shared', ignore_errors=True) \n    return pl.DataFrame(data)\n\n\n# ----------------------- RUN KAGGLE INFERENCE SERVER & FINAL SUBMISSION ------------------------\n\nprint(\"\\nStarting Kaggle Evaluation API Interface...\")\nshutil.rmtree('/kaggle/shared', ignore_errors=True)\n\n# Initialize the inference server with our custom predict function\ninference_server = kaggle_evaluation.rsna_inference_server.RSNAInferenceServer(predict)\n\nif os.getenv('KAGGLE_IS_COMPETITION_RERUN'):\n    # OFFICIAL SUBMISSION MODE\n    print(\"Running in official submission RERUN mode.\")\n    inference_server.serve()\nelse:\n    # LOCAL TEST MODE\n    print(\"Running in local TEST mode (generating submission.parquet)...\")\n    inference_server.run_local_gateway()\n    try:\n        display(pl.read_parquet('/kaggle/working/submission.parquet'))\n    except Exception as e:\n        print(f\"Could not display submission.parquet: {e}\")\n\nprint(\"✅ Inference and submission process finished.\")","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null}]}