{"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.12"},"kaggle":{"accelerator":"nvidiaTeslaT4","dataSources":[{"sourceId":92399,"databundleVersionId":11038207,"isSourceIdPinned":false,"sourceType":"competition"},{"sourceId":10935081,"sourceType":"datasetVersion","datasetId":6799927}],"dockerImageVersionId":30918,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":true},"papermill":{"default_parameters":{},"duration":8987.706787,"end_time":"2025-04-12T06:36:02.594808","environment_variables":{},"exception":null,"input_path":"__notebook__.ipynb","output_path":"__notebook__.ipynb","parameters":{},"start_time":"2025-04-12T04:06:14.888021","version":"2.6.0"}},"nbformat_minor":4,"nbformat":4,"cells":[{"id":"30899270","cell_type":"markdown","source":"# Dashcam Collision Prediction Project 🚗\n\n## Overview 📜\nThis project implements a deep learning pipeline for the Nexar Collision Prediction competition, which aims to predict potential vehicle collisions from dashcam footage. Early detection of imminent collisions is crucial for advanced driver assistance systems and autonomous driving technology to prevent accidents and improve road safety. 🛣️\n\n## Key Components 🧩\n\n### 1. Data Preprocessing 📊\n* Efficient keyframe extraction focusing on critical moments in videos\n* Computation of optical flow to capture motion dynamics\n* Augmentation techniques simulating various driving conditions\n* Parallel processing with multiprocessing Pool for scalability\n* Handling of video artifacts and missing frames\n\n### 2. Feature Engineering and Extraction\n* Adaptive frame sampling with emphasis on final seconds where collisions occur\n* Motion representation through Farneback optical flow algorithm\n* Realistic data augmentation with random flips, color jitter, fog, and rain effects\n* Custom transformations preserving temporal coherence across video frames\n* Efficient NumPy to PyTorch conversion maintaining batch dimensionality\n\n### 3. Multi-Stream Architecture 🧮\n* Visual stream based on MobileNetV2 for spatial feature extraction\n* Specialized optical flow stream using 3D convolutions for motion analysis\n* Temporal modeling with LSTM for sequence dynamics\n* Lightweight implementation optimized for memory efficiency\n* Feature fusion combining spatial and temporal information\n\n### 4. Model Optimization 🎨\n* Memory-efficient processing for video batches\n* GPU-compatible training with automatic fallback to CPU\n* Gradient clipping to ensure training stability\n* Learning rate scheduling with ReduceLROnPlateau\n* Checkpoint management to save best models\n\n### 5. Custom Loss Function\n* Combined binary classification and regression objectives\n* Specialized time prediction loss with late alert penalty\n* Balanced weighting between collision detection and timing\n* Handling of edge cases and outliers in training data\n* Focus on practical safety impacts in loss design\n\n### 6. Evaluation and Inference 📏🌌\n* Dynamic thresholding based on video conditions\n* Uncertainty quantification with Monte Carlo dropout\n* Memory-efficient batch prediction for large test sets\n* Visualization capabilities for model interpretation\n* Submission format generation for competition evaluation\n\n## Methodology 🔍\nThe pipeline adopts a multi-stream approach that processes both spatial and temporal information from dashcam videos. It integrates computer vision techniques with deep learning to create a system capable of early collision warning.\n\nThe key aspects of the methodology include:\n* Two-stream architecture separating appearance and motion analysis\n* Temporal modeling to capture the dynamics of traffic situations\n* Custom data augmentation mimicking real-world driving conditions\n* Memory-optimized implementation suitable for video processing\n* Balanced training approach handling both classification and timing tasks\n\nThis lightweight model delivers effective collision prediction while maintaining computational efficiency, making it suitable for potential deployment in actual vehicle safety systems.","metadata":{"papermill":{"duration":0.005925,"end_time":"2025-04-12T04:06:17.675363","exception":false,"start_time":"2025-04-12T04:06:17.669438","status":"completed"},"tags":[]}},{"id":"d5e21c20","cell_type":"markdown","source":"## Import necessary libraries","metadata":{"papermill":{"duration":0.004675,"end_time":"2025-04-12T04:06:17.685275","exception":false,"start_time":"2025-04-12T04:06:17.680600","status":"completed"},"tags":[]}},{"id":"6b89623b","cell_type":"code","source":"import os\nimport gc\nimport time\nimport warnings\nfrom multiprocessing import Pool\n\nimport cv2\nimport numpy as np\nimport pandas as pd\nfrom sklearn.metrics import roc_auc_score\nfrom sklearn.model_selection import KFold\nfrom sklearn.metrics import confusion_matrix, precision_recall_curve\nimport matplotlib.pyplot as plt\nimport seaborn as sns\nimport torch\nimport torch.nn as nn\nimport torch.nn.functional as F\nimport torch.optim as optim\nfrom torch.utils.data import Dataset, DataLoader\nimport torchvision\nimport torchvision.models\nfrom torch.optim.lr_scheduler import CosineAnnealingWarmRestarts\nfrom torch.optim.lr_scheduler import CyclicLR\ntorch.nn.utils.clip_grad_norm_\n\nwarnings.filterwarnings(\"ignore\")\n\n# Check GPU availability and set device\ndevice = 'cuda' if torch.cuda.is_available() else 'cpu'\nprint(f\"Using device: {device}\")","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:17.569804Z","iopub.execute_input":"2025-04-12T23:35:17.570114Z","iopub.status.idle":"2025-04-12T23:35:28.694373Z","shell.execute_reply.started":"2025-04-12T23:35:17.570090Z","shell.execute_reply":"2025-04-12T23:35:28.693560Z"},"papermill":{"duration":11.046289,"end_time":"2025-04-12T04:06:28.736694","exception":false,"start_time":"2025-04-12T04:06:17.690405","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"85117a1a","cell_type":"code","source":"# Suppress unnecessary formatting warnings\nwarnings.filterwarnings(\"ignore\", category=RuntimeWarning)\n\n# Paths to the CSV files\ntrain_csv_path = '/kaggle/input/nexar-collision-prediction/train.csv'\ntest_csv_path = '/kaggle/input/nexar-collision-prediction/test.csv'\nsubmission_csv_path = '/kaggle/input/nexar-collision-prediction/sample_submission.csv'\n\n# Load the CSV files\ntrain_df = pd.read_csv(train_csv_path)\ntest_df = pd.read_csv(test_csv_path)\nsubmission_df = pd.read_csv(submission_csv_path)\n\n# Display the first few rows of the DataFrames\nprint(\"Train.csv:\")\nprint(train_df.head())\n\nprint(\"\\nTest.csv:\")\nprint(test_df.head())\n\nprint(\"\\nSample Submission:\")\nprint(submission_df.head())\n\n# Optional: handle NaN values if needed, filling with zero or another value\ntrain_df['time_of_event'] = train_df['time_of_event'].fillna(0)\ntrain_df['time_of_alert'] = train_df['time_of_alert'].fillna(0)","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.695559Z","iopub.execute_input":"2025-04-12T23:35:28.696032Z","iopub.status.idle":"2025-04-12T23:35:28.755119Z","shell.execute_reply.started":"2025-04-12T23:35:28.696005Z","shell.execute_reply":"2025-04-12T23:35:28.754474Z"},"papermill":{"duration":0.052565,"end_time":"2025-04-12T04:06:28.794827","exception":false,"start_time":"2025-04-12T04:06:28.742262","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"7f075baa","cell_type":"markdown","source":"## Data Preprocessing and Feature Extraction","metadata":{"papermill":{"duration":0.005383,"end_time":"2025-04-12T04:06:28.805129","exception":false,"start_time":"2025-04-12T04:06:28.799746","status":"completed"},"tags":[]}},{"id":"3fc148ba","cell_type":"code","source":"def extract_keyframes(video_path, num_frames=12, target_size=(160, 160)):\n    \"\"\"\n    Extracts key frames from the video, focusing on the final part where collisions typically occur.\n    Uses exponential distribution to give more weight to frames closer to the end.\n    \"\"\"\n    cap = cv2.VideoCapture(video_path)\n    \n    if not cap.isOpened():\n        print(f\"Could not open the video: {video_path}\")\n        return np.zeros((num_frames, target_size[0], target_size[1], 3), dtype=np.uint8)\n    \n    frames = []\n    total_frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT))\n    fps = cap.get(cv2.CAP_PROP_FPS)\n    \n    if total_frames <= 0:\n        print(f\"Video without frames: {video_path}\")\n        cap.release()\n        return np.zeros((num_frames, target_size[0], target_size[1], 3), dtype=np.uint8)\n    \n    # Calculate video duration in seconds\n    duration = total_frames / fps if fps > 0 else 0\n    \n    # If the video is short (less than 10 seconds), distribute frames uniformly\n    if duration < 10:\n        frame_indices = np.linspace(0, total_frames - 1, num_frames, dtype=int)\n    else:\n        # Concentrate 80% of frames in the last 3 seconds (critical area)\n        end_frames = int(num_frames * 0.8)\n        start_frames = num_frames - end_frames\n        \n        # Calculate the starting index for the last 3 seconds\n        last_seconds = 3\n        last_frame_count = min(int(fps * last_seconds), total_frames - 1)\n        start_idx = max(0, total_frames - last_frame_count)\n        \n        # Exponential distribution to give more weight to the last frames\n        # This creates indices that are more densely packed toward the end\n        end_indices = np.array([\n            start_idx + int((total_frames - start_idx - 1) * (i/end_frames)**2) \n            for i in range(1, end_frames + 1)\n        ])\n        \n        # Initial frames distributed uniformly for context\n        begin_indices = np.linspace(0, start_idx - 1, start_frames, dtype=int) if start_idx > 0 else np.zeros(start_frames, dtype=int)\n        \n        # Combine indices\n        frame_indices = np.concatenate([begin_indices, end_indices])\n    \n    # Extract selected frames\n    for idx in frame_indices:\n        cap.set(cv2.CAP_PROP_POS_FRAMES, idx)\n        ret, frame = cap.read()\n        if ret:\n            # Use higher resolution and better interpolation\n            frame = cv2.resize(frame, target_size, interpolation=cv2.INTER_LANCZOS4)\n            frame = cv2.cvtColor(frame, cv2.COLOR_BGR2RGB)\n            frames.append(frame)\n        else:\n            frames.append(np.zeros((target_size[0], target_size[1], 3), dtype=np.uint8))\n    \n    cap.release()\n    return np.array(frames, dtype=np.uint8)\n\n# First, define transformation classes in the global scope\nclass RandomHorizontalFlip(object):\n    def __init__(self, p=0.5):\n        self.p = p\n        \n    def __call__(self, frames):\n        if np.random.random() < self.p:\n            return frames[:, :, ::-1, :].copy()  # horizontally flip each frame\n        return frames\n\nclass ColorJitter(object):\n    def __init__(self, brightness=0, contrast=0):\n        self.brightness = brightness\n        self.contrast = contrast\n        \n    def __call__(self, frames):\n        # Apply brightness jitter\n        if self.brightness > 0:\n            brightness_factor = np.random.uniform(max(0, 1-self.brightness), 1+self.brightness)\n            frames = frames * brightness_factor\n            frames = np.clip(frames, 0, 255)\n        \n        # Apply contrast jitter\n        if self.contrast > 0:\n            contrast_factor = np.random.uniform(max(0, 1-self.contrast), 1+self.contrast)\n            frames = (frames - 128) * contrast_factor + 128\n            frames = np.clip(frames, 0, 255)\n            \n        return frames\n\nclass AddFog(object):\n    def __call__(self, frames):\n        fog = np.random.uniform(0.7, 0.9, frames.shape).astype(np.float32)\n        return frames * 0.8 + fog * 50  # Adjusted for 0-255 scale\n\nclass AddRain(object):\n    def __call__(self, frames):\n        h, w = frames.shape[1:3]\n        rain = np.random.uniform(0, 1, (len(frames), h, w, 1)).astype(np.float32)\n        rain = (rain > 0.97).astype(np.float32) * 200  # White rain drops\n        return np.clip(frames * 0.9 + rain, 0, 255)  # Darken a bit and add drops\n\nclass RandomApply(object):\n    def __init__(self, transform, p=0.5):\n        self.transform = transform\n        self.p = p\n        \n    def __call__(self, frames):\n        if np.random.random() < self.p:\n            return self.transform(frames)\n        return frames\n\nclass Compose(object):\n    def __init__(self, transforms):\n        self.transforms = transforms\n        \n    def __call__(self, frames):\n        for t in self.transforms:\n            frames = t(frames)\n        return frames\n\nclass ToTensor(object):\n    def __call__(self, frames):\n        # Convert from (T, H, W, C) to (T, C, H, W)\n        frames = frames.transpose(0, 3, 1, 2)\n        # Convert to tensor and normalize to [0, 1]\n        return torch.from_numpy(frames).float() / 255.0\n\ndef get_video_transforms():\n    \"\"\"\n    Returns transformations for data augmentation in videos.\n    \"\"\"\n    return {\n        'train': Compose([\n            RandomHorizontalFlip(p=0.5),\n            ColorJitter(brightness=0.3, contrast=0.3),\n            RandomApply(AddFog(), p=0.15),\n            RandomApply(AddRain(), p=0.15),\n            RandomApply(RandomNoise(0.05), p=0.2), \n            RandomApply(RandomOcclusion(), p=0.1),\n            # Adicione estas novas transformações\n            RandomApply(AddSunGlare(), p=0.1),\n            RandomApply(SimulateNight(), p=0.1),\n            RandomApply(Motion(), p=0.2),\n            ToTensor()\n        ]),\n        'val': Compose([\n            ToTensor()  # Only tensor conversion for validation\n        ])\n    }\n\nclass AddSunGlare(object):\n    \"\"\"Simulates the effect of sun glare on the windshield.\"\"\"\n    def __call__(self, frames):\n        h, w = frames.shape[1:3]\n        \n        # Random position for sun glare\n        x = np.random.randint(0, w)\n        y = np.random.randint(0, h // 3)  # More likely at the top\n        radius = np.random.randint(30, 80)\n        intensity = np.random.uniform(0.6, 0.9)\n        \n        result = frames.copy().astype(np.float32)\n        \n        for i in range(len(frames)):\n            # Create a circular gradient to simulate glare\n            Y, X = np.ogrid[:h, :w]\n            dist = np.sqrt((X - x) ** 2 + (Y - y) ** 2)\n            mask = dist <= radius\n            \n            # Apply glare with gradient\n            glow = np.maximum(0, (1 - dist / radius)) * intensity * 255\n            result[i, :, :, :] += np.repeat(glow[:, :, np.newaxis], 3, axis=2) * mask[:, :, np.newaxis]\n            \n        return np.clip(result, 0, 255).astype(np.uint8)\n\nclass SimulateNight(object):\n    \"\"\"Simulates low-light/night conditions.\"\"\"\n    def __call__(self, frames):\n        # Reduce overall brightness\n        darkness = np.random.uniform(0.3, 0.6)\n        result = frames.astype(np.float32) * darkness\n        \n        # Add noise to simulate low-light sensor conditions\n        noise = np.random.normal(0, 5, frames.shape).astype(np.float32)\n        result += noise\n        \n        return np.clip(result, 0, 255).astype(np.uint8)\n\nclass Motion(object):\n    \"\"\"Simulates motion blur.\"\"\"\n    def __call__(self, frames):\n        # Choose direction and magnitude of motion blur\n        kernel_size = np.random.choice([3, 5, 7])\n        angle = np.random.uniform(0, 360)\n        \n        # Create motion blur kernel\n        kernel = np.zeros((kernel_size, kernel_size))\n        center = kernel_size // 2\n        \n        # Draw a line at the specified angle\n        x1 = center\n        y1 = center\n        x2 = int(center + (kernel_size - 1) / 2 * np.cos(np.radians(angle)))\n        y2 = int(center + (kernel_size - 1) / 2 * np.sin(np.radians(angle)))\n        \n        cv2.line(kernel, (x1, y1), (x2, y2), 1.0)\n        kernel = kernel / np.sum(kernel)\n        \n        result = np.zeros_like(frames)\n        \n        # Apply the kernel to each frame and channel\n        for i in range(len(frames)):\n            for c in range(3):\n                result[i, :, :, c] = cv2.filter2D(frames[i, :, :, c], -1, kernel)\n                \n        return result\n\nclass RandomNoise(object):\n    \"\"\"\n    Applies random Gaussian noise to video frames for data augmentation.\n    \n    This transformation helps the model become more robust to noise\n    that may be present in real-world video data.\n    \n    Args:\n        std (float): Standard deviation of the Gaussian noise as a fraction\n                     of the pixel value range (default: 0.05)\n    \"\"\"\n    def __init__(self, std=0.05):\n        self.std = std\n        \n    def __call__(self, frames):\n        \"\"\"\n        Apply random noise to the input frames.\n        \n        Args:\n            frames (numpy.ndarray): Input video frames of shape (T, H, W, C)\n                                   where T is number of frames\n        \n        Returns:\n            numpy.ndarray: Noise-augmented frames, clipped to valid pixel range [0, 255]\n        \"\"\"\n        # Generate Gaussian noise with specified standard deviation\n        noise = np.random.normal(0, self.std * 255, frames.shape).astype(np.float32)\n        \n        # Add noise and clip to valid pixel range\n        return np.clip(frames + noise, 0, 255).astype(np.uint8)\n\n\nclass RandomOcclusion(object):\n    \"\"\"\n    Simulates occlusion in video frames by adding black rectangles.\n    \n    This transformation helps the model learn to handle partial occlusions\n    that may occur in real-world scenarios when objects block the camera view.\n    \"\"\"\n    def __call__(self, frames):\n        \"\"\"\n        Apply random occlusion to the input frames.\n        \n        Args:\n            frames (numpy.ndarray): Input video frames of shape (T, H, W, C)\n                                   where T is number of frames\n        \n        Returns:\n            numpy.ndarray: Frames with random occlusion applied\n        \"\"\"\n        # Get frame dimensions\n        h, w = frames.shape[1:3]\n        \n        # Define occlusion area: 10-25% of the image\n        occl_h = np.random.randint(int(h * 0.1), int(h * 0.25))\n        occl_w = np.random.randint(int(w * 0.1), int(w * 0.25))\n        \n        # Randomly position the occlusion\n        occl_x = np.random.randint(0, w - occl_w)\n        occl_y = np.random.randint(0, h - occl_h)\n        \n        # Create a copy to avoid modifying the original frames\n        frames_copy = frames.copy()\n        \n        # Apply occlusion to all frames by setting pixels to zero (black)\n        for i in range(len(frames)):\n            frames_copy[i, occl_y:occl_y+occl_h, occl_x:occl_x+occl_w, :] = 0\n            \n        return frames_copy\n        \ndef compute_optical_flow(frames, skip_frames=1):\n    \"\"\"Calculates optical flow skipping some frames to reduce processing.\"\"\"\n    if len(frames) < 2:\n        return np.zeros((1, frames.shape[1], frames.shape[2], 2), dtype=np.float32)\n    \n    flows = []\n    prev_gray = cv2.cvtColor(frames[0], cv2.COLOR_RGB2GRAY)\n    \n    for i in range(1, len(frames), skip_frames):\n        curr_gray = cv2.cvtColor(frames[i], cv2.COLOR_RGB2GRAY)\n        try:\n            # Reduce parameters for faster calculation\n            flow = cv2.calcOpticalFlowFarneback(prev_gray, curr_gray,\n                                               None, 0.5, 3, 15, 3, 5, 1.2, 0)\n            flows.append(flow)\n        except Exception as e:\n            print(f\"Error calculating optical flow: {str(e)}\")\n            flows.append(np.zeros((frames.shape[1], frames.shape[2], 2), dtype=np.float32))\n            \n        prev_gray = curr_gray\n    \n    if not flows:\n        return np.zeros((1, frames.shape[1], frames.shape[2], 2), dtype=np.float32)\n        \n    return np.array(flows, dtype=np.float32)\n\ndef process_video(args):\n    \"\"\"\n    Function to process an individual video.\n    \"\"\"\n    video_path, video_id, num_frames = args\n    try:\n        # Extract frames with higher resolution\n        frames = extract_keyframes(video_path, num_frames=num_frames, target_size=(160, 160))\n        \n        # Calculate optical flow\n        optical_flow = compute_optical_flow(frames, skip_frames=1)\n        \n        # We return NumPy arrays instead of applying transformations now\n        return video_id, {\n            'frames': frames,\n            'optical_flow': optical_flow,\n        }\n    except Exception as e:\n        print(f\"Error processing video {video_id}: {str(e)}\")\n        return video_id, None\n\ndef parallel_preprocess_dataset(video_dir, video_ids, num_frames=8, num_workers=4):\n    \"\"\"\n    Pre-processes multiple videos in parallel.\n    \"\"\"\n    args_list = []\n    for video_id in video_ids:\n        video_path = os.path.join(video_dir, f\"{video_id}.mp4\")\n        if os.path.exists(video_path):\n            args_list.append((video_path, video_id, num_frames))\n    \n    start_time = time.time()\n    print(f\"Starting parallel pre-processing of {len(args_list)} videos with {num_workers} workers...\")\n    \n    processed_data = {}\n    with Pool(num_workers) as p:\n        results = p.map(process_video, args_list)\n        for video_id, data in results:\n            if data is not None:\n                processed_data[video_id] = data\n    \n    print(f\"Pre-processing completed in {time.time() - start_time:.2f} seconds.\")\n    print(f\"Processed {len(processed_data)} out of {len(args_list)} videos.\")\n    \n    return processed_data","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.756520Z","iopub.execute_input":"2025-04-12T23:35:28.756736Z","iopub.status.idle":"2025-04-12T23:35:28.789572Z","shell.execute_reply.started":"2025-04-12T23:35:28.756717Z","shell.execute_reply":"2025-04-12T23:35:28.788839Z"},"papermill":{"duration":0.039068,"end_time":"2025-04-12T04:06:28.849051","exception":false,"start_time":"2025-04-12T04:06:28.809983","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"8336bc24","cell_type":"markdown","source":"## Custom Dataset Creation","metadata":{"papermill":{"duration":0.004257,"end_time":"2025-04-12T04:06:28.858234","exception":false,"start_time":"2025-04-12T04:06:28.853977","status":"completed"},"tags":[]}},{"id":"feb5c68e","cell_type":"code","source":"class DashcamDataset(Dataset):\n   def __init__(self, video_dir, annotations, transform=None, num_frames=8):\n       self.video_dir = video_dir\n       self.annotations = annotations\n       self.transform = transform\n       self.num_frames = num_frames\n       self.video_ids = list(annotations.keys())\n\n   def __len__(self):\n       return len(self.video_ids)\n\n   def __getitem__(self, idx):\n       video_id = self.video_ids[idx]\n       video_path = os.path.join(self.video_dir, f\"{video_id}.mp4\")\n       \n       try:\n           # Check if the file exists\n           if not os.path.exists(video_path):\n               print(f\"Video not found: {video_path}\")\n               raise FileNotFoundError(f\"File not found: {video_path}\")\n               \n           # Extract frames with reduced resolution\n           frames = extract_keyframes(video_path, self.num_frames, target_size=(112, 112))\n           \n           # Calculate optical flow with skip_frames\n           optical_flow = compute_optical_flow(frames, skip_frames=1)\n           \n           # Apply transformations\n           if self.transform:\n               frames = self.transform(frames)\n           else:\n               # Convert to tensor manually\n               frames = torch.from_numpy(frames.transpose(0, 3, 1, 2)).float() / 255.0\n           \n           # Convert optical flow to tensor\n           optical_flow = torch.from_numpy(optical_flow.transpose(0, 3, 1, 2)).float()\n           \n           # Load label and alert time\n           label = self.annotations[video_id]['label']\n           alert_time = self.annotations[video_id].get('alert_time', 0)\n           \n           return {\n               'frames': frames,\n               'optical_flow': optical_flow,\n               'label': torch.tensor(label).float(),\n               'alert_time': torch.tensor(alert_time).float(),\n               'video_id': video_id\n           }\n               \n       except Exception as e:\n           print(f\"Error processing video {video_id}: {str(e)}\")\n           # Create a placeholder for this video\n           dummy_frames = torch.zeros((self.num_frames, 3, 112, 112))\n           dummy_flow = torch.zeros((max(1, self.num_frames-1), 2, 112, 112))\n           return {\n               'frames': dummy_frames,\n               'optical_flow': dummy_flow,\n               'label': torch.tensor(self.annotations[video_id]['label']).float(),\n               'alert_time': torch.tensor(self.annotations[video_id].get('alert_time', 0)).float(),\n               'video_id': video_id\n           }\n\nclass PreprocessedDashcamDataset(Dataset):\n   def __init__(self, processed_data, annotations, transform=None):\n       self.processed_data = processed_data\n       self.annotations = annotations\n       self.transform = transform\n       self.video_ids = list(processed_data.keys())\n       \n   def __len__(self):\n       return len(self.video_ids)\n   \n   def __getitem__(self, idx):\n       video_id = self.video_ids[idx]\n       data = self.processed_data[video_id]\n       \n       # Get frames and optical flow\n       frames = data['frames']\n       optical_flow = data['optical_flow']\n       \n       # Apply transformations to frames\n       if self.transform:\n           frames_tensor = self.transform(frames)\n       else:\n           # Convert to tensor manually\n           frames_tensor = torch.from_numpy(frames.transpose(0, 3, 1, 2)).float() / 255.0\n       \n       # Convert optical flow to tensor\n       optical_flow_tensor = torch.from_numpy(optical_flow.transpose(0, 3, 1, 2)).float()\n       \n       # Load label and alert time\n       label = self.annotations[video_id]['label']\n       alert_time = self.annotations[video_id].get('alert_time', 0)\n       \n       return {\n           'frames': frames_tensor,\n           'optical_flow': optical_flow_tensor,\n           'label': torch.tensor(label).float(),\n           'alert_time': torch.tensor(alert_time).float(),\n           'video_id': video_id\n       }","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.790541Z","iopub.execute_input":"2025-04-12T23:35:28.790813Z","iopub.status.idle":"2025-04-12T23:35:28.815958Z","shell.execute_reply.started":"2025-04-12T23:35:28.790791Z","shell.execute_reply":"2025-04-12T23:35:28.815294Z"},"papermill":{"duration":0.016911,"end_time":"2025-04-12T04:06:28.879816","exception":false,"start_time":"2025-04-12T04:06:28.862905","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"39617cbf","cell_type":"markdown","source":"## Definition of the Multi-Stream Model with Fusion via Transformer","metadata":{"papermill":{"duration":0.004316,"end_time":"2025-04-12T04:06:28.888570","exception":false,"start_time":"2025-04-12T04:06:28.884254","status":"completed"},"tags":[]}},{"id":"b88c69d2","cell_type":"code","source":"# Spatial Attention mechanism\nclass SpatialAttention(nn.Module):\n    def __init__(self, in_channels):\n        super(SpatialAttention, self).__init__()\n        self.conv = nn.Sequential(\n            nn.Conv2d(in_channels, in_channels // 8, kernel_size=1),\n            nn.ReLU(),\n            nn.Conv2d(in_channels // 8, 1, kernel_size=1),\n            nn.Sigmoid()\n        )\n    \n    def forward(self, x):\n        # x: [batch, channels, height, width]\n        attention_map = self.conv(x)\n        return x * attention_map\n\nclass MobileNetVisualStream(nn.Module):\n    def __init__(self, output_features=64):\n        super(MobileNetVisualStream, self).__init__()\n        \n        print(\"Initializing MobileNetV2 model...\")\n        \n        # Create the model WITHOUT pre-trained weights first\n        mobilenet = torchvision.models.mobilenet_v2(pretrained=False)\n        \n        # Check if we have local Hugging Face weights\n        hf_weights_path = \"/kaggle/input/googlemobilenet-v2-1-0-224/pytorch_model.bin\"\n        \n        if os.path.exists(hf_weights_path):\n            print(f\"Found Hugging Face pre-trained weights: {hf_weights_path}\")\n            try:\n                # Load the weights from Hugging Face\n                state_dict = torch.load(hf_weights_path, map_location='cpu')\n                \n                # Map the weights from the Hugging Face format to the torchvision format\n                mobilenet_state_dict = {}\n                for k, v in state_dict.items():\n                    if k.startswith('mobilenet_v2.'):\n                        # Remove the prefix 'mobilenet_v2.'\n                        mobilenet_state_dict[k[13:]] = v\n                \n                # Load the mapped weights\n                mobilenet.load_state_dict(mobilenet_state_dict, strict=False)\n                print(\"Successfully loaded weights from Hugging Face model\")\n            except Exception as e:\n                print(f\"Error loading Hugging Face weights: {e}\")\n                print(\"Using random initialization\")\n        else:\n            print(\"Hugging Face weights not found. Using random initialization.\")\n        \n        # Divide into blocks for adding attention between them\n        self.block1 = nn.Sequential(*list(mobilenet.features)[:7])  # First layers\n        self.spatial_attention1 = SpatialAttention(32)  # Adjust for number of channels\n        \n        self.block2 = nn.Sequential(*list(mobilenet.features)[7:14])  # Middle layers\n        self.spatial_attention2 = SpatialAttention(96)  # Adjust for number of channels\n        \n        self.block3 = nn.Sequential(*list(mobilenet.features)[14:])  # Last layers\n        \n        # Projection layer to reduce dimensionality\n        self.projection = nn.Sequential(\n            nn.AdaptiveAvgPool2d(1),\n            nn.Flatten(),\n            nn.Linear(1280, output_features),\n            nn.ReLU()\n        )\n        \n        self.output_features = output_features\n        \n    def forward(self, x):\n        # x: [batch, T, C, H, W]\n        batch_size, seq_len, c, h, w = x.size()\n        \n        # Process each frame individually\n        features = []\n        for t in range(seq_len):\n            # Extract features with spatial attention\n            feat = self.block1(x[:, t])\n            feat = self.spatial_attention1(feat)\n            \n            feat = self.block2(feat)\n            feat = self.spatial_attention2(feat)\n            \n            feat = self.block3(feat)\n            feat = self.projection(feat)\n            \n            features.append(feat)\n        \n        # Stack along the temporal dimension\n        output = torch.stack(features, dim=1)  # [batch, T, output_features]\n        \n        return output\n\nclass TemporalAttention(nn.Module):\n    def __init__(self, dim):\n        super(TemporalAttention, self).__init__()\n        self.query = nn.Linear(dim, dim)  # Query projection\n        self.key = nn.Linear(dim, dim)    # Key projection\n        self.value = nn.Linear(dim, dim)  # Value projection\n        self.scale = dim ** -0.5          # Scaling factor for dot products\n        \n    def forward(self, x):\n        # Project inputs to queries, keys and values\n        q = self.query(x)  # [batch, seq_len, dim]\n        k = self.key(x)    # [batch, seq_len, dim]\n        v = self.value(x)  # [batch, seq_len, dim]\n        \n        # Compute scaled dot-product attention\n        attn = (q @ k.transpose(-2, -1)) * self.scale  # [batch, seq_len, seq_len]\n        attn = F.softmax(attn, dim=-1)                 # Apply softmax to get attention weights\n        \n        # Apply attention weights to values\n        out = attn @ v  # [batch, seq_len, dim]\n        return out\n\n# Optimized version of OpticalFlowStream class\nclass OpticalFlowStream(nn.Module):\n    def __init__(self, output_features=48):\n        super(OpticalFlowStream, self).__init__()\n        self.features = nn.Sequential(\n            nn.Conv3d(2, 24, kernel_size=(3, 3, 3), padding=1),\n            nn.BatchNorm3d(24),\n            nn.ReLU(),\n            nn.MaxPool3d(kernel_size=(1, 2, 2)),\n\n            nn.Conv3d(24, 48, kernel_size=(3, 3, 3), padding=1),\n            nn.BatchNorm3d(48),\n            nn.ReLU(),\n            nn.MaxPool3d(kernel_size=(1, 2, 2)),\n\n            nn.Conv3d(48, output_features, kernel_size=(3, 3, 3), padding=1),\n            nn.BatchNorm3d(output_features),\n            nn.ReLU(),\n            nn.AdaptiveAvgPool3d((None, 1, 1))\n        )\n        \n        # Add temporal attention specific to optical flow\n        self.temporal_attention = nn.Sequential(\n            nn.Conv1d(output_features, output_features, kernel_size=3, padding=1),\n            nn.Sigmoid()\n        )\n    \n    def forward(self, x):\n        # x: [batch, T, C, H, W] -> [batch, C, T, H, W]\n        x = x.permute(0, 2, 1, 3, 4)\n        x = self.features(x)\n        features = x.squeeze(-1).squeeze(-1)  # [batch, features, T]\n        \n        # Apply temporal attention\n        attention = self.temporal_attention(features)\n        weighted_features = features * attention\n        \n        return weighted_features\n\n# Optimized version of FusionTransformer class\nclass FusionTransformer(nn.Module):\n   def __init__(self, visual_dim=64, flow_dim=32, output_dim=64, nhead=4, num_layers=2):\n       super(FusionTransformer, self).__init__()\n       # Use simpler architecture to save memory\n       self.visual_proj = nn.Linear(visual_dim, output_dim)\n       self.flow_proj = nn.Linear(flow_dim, output_dim)\n       \n       # Use an LSTM layer instead of Transformer (less memory intensive)\n       self.lstm = nn.LSTM(\n           input_size=output_dim,\n           hidden_size=output_dim,\n           num_layers=1,\n           batch_first=True\n       )\n       \n       # Output layer\n       self.output_layer = nn.Linear(output_dim, output_dim)\n   \n   def forward(self, visual_feat, flow_feat):\n       # Temporal pooling\n       visual_pool = visual_feat.mean(dim=2)  # [batch, visual_dim]\n       flow_pool = flow_feat.mean(dim=2)      # [batch, flow_dim]\n       \n       # Projection to common dimension\n       visual_proj = self.visual_proj(visual_pool)  # [batch, output_dim]\n       flow_proj = self.flow_proj(flow_pool)        # [batch, output_dim]\n       \n       # Concatenate as sequence\n       concat = torch.stack([visual_proj, flow_proj], dim=1)  # [batch, 2, output_dim]\n       \n       # Pass through LSTM\n       lstm_out, _ = self.lstm(concat)\n       \n       # Final pooling\n       output = lstm_out.mean(dim=1)        # [batch, output_dim]\n       output = self.output_layer(output)   # [batch, output_dim]\n       \n       return output\n\nclass TemporalAttention(nn.Module):\n    \"\"\"\n    Temporal attention mechanism that learns to focus on important time steps in a sequence.\n    \n    This module implements a self-attention mechanism where each time step attends\n    to all other time steps, allowing the model to capture temporal dependencies.\n    \n    Args:\n        dim (int): Dimension of the input feature space\n    \"\"\"\n    def __init__(self, dim):\n        super(TemporalAttention, self).__init__()\n        self.query = nn.Linear(dim, dim)  # Query projection\n        self.key = nn.Linear(dim, dim)    # Key projection\n        self.value = nn.Linear(dim, dim)  # Value projection\n        self.scale = dim ** -0.5          # Scaling factor for dot products\n        \n    def forward(self, x):\n        \"\"\"\n        Forward pass of the temporal attention mechanism.\n        \n        Args:\n            x (torch.Tensor): Input tensor of shape [batch, seq_len, dim]\n            \n        Returns:\n            torch.Tensor: Attention-weighted output of same shape as input\n        \"\"\"\n        # Project inputs to queries, keys and values\n        q = self.query(x)  # [batch, seq_len, dim]\n        k = self.key(x)    # [batch, seq_len, dim]\n        v = self.value(x)  # [batch, seq_len, dim]\n        \n        # Compute scaled dot-product attention\n        attn = (q @ k.transpose(-2, -1)) * self.scale  # [batch, seq_len, seq_len]\n        attn = F.softmax(attn, dim=-1)                 # Apply softmax to get attention weights\n        \n        # Apply attention weights to values\n        out = attn @ v  # [batch, seq_len, dim]\n        return out\n\n\nclass LightweightMultiStreamModel(nn.Module):\n    def __init__(self, output_dim=64):\n        super(LightweightMultiStreamModel, self).__init__()\n        # Optimized parameters\n        visual_features = 64\n        flow_features = 32\n        \n        # Visual and flow processing streams\n        self.visual_stream = MobileNetVisualStream(output_features=visual_features)\n        self.flow_stream = OpticalFlowStream(output_features=flow_features)\n        \n        # Projection layers to align feature dimensions\n        self.visual_proj = nn.Linear(visual_features, output_dim)\n        self.flow_proj = nn.Linear(flow_features, output_dim)\n        \n        # Temporal attention mechanism\n        self.temporal_attention = TemporalAttention(output_dim)\n        \n        # Fusion layer after attention\n        self.fusion_layer = nn.Sequential(\n            nn.Linear(output_dim, output_dim),\n            nn.ReLU()\n        )\n        \n        # Stronger dropout for regularization\n        self.dropout = nn.Dropout(0.5)\n        \n        # Classification branch with additional regularization\n        self.classifier = nn.Sequential(\n            nn.Linear(output_dim, 32),\n            nn.ReLU(),\n            nn.Dropout(0.4),\n            nn.Linear(32, 1),\n            nn.Sigmoid()\n        )\n        \n        # Regression branch for alert time prediction\n        self.regressor = nn.Sequential(\n            nn.Linear(output_dim, 32),\n            nn.ReLU(),\n            nn.Dropout(0.3),\n            nn.Linear(32, 1)\n        )\n        \n        # Flag to track if we've created dynamic projections already\n        self.dynamic_projections_created = False\n    \n    def forward(self, frames, optical_flow):\n        # Process visual and optical flow streams\n        visual_feat = self.visual_stream(frames)  # [batch, T, features]\n        flow_feat = self.flow_stream(optical_flow)  # [batch, features, T-1]\n        \n        # Check tensor shapes for debugging\n        batch_size = frames.size(0)\n    \n        # Reshape visual features based on their dimensions\n        if len(visual_feat.shape) == 5:  # If [batch, features, 1, 1, 1]\n            visual_feat = visual_feat.squeeze(-1).squeeze(-1).squeeze(-1)  # [batch, features]\n            visual_feat = visual_feat.unsqueeze(1)  # [batch, 1, features]\n        elif len(visual_feat.shape) == 3:  # If [batch, features, T]\n            visual_feat = visual_feat.permute(0, 2, 1)  # [batch, T, features]\n        elif len(visual_feat.shape) == 2:  # If [batch, T*features]\n            # Assuming num_frames is a known constant (e.g., 8)\n            num_frames = frames.size(1)\n            visual_feat = visual_feat.view(batch_size, num_frames, -1)  # [batch, T, features_per_frame]\n    \n        # Reshape optical flow features based on their dimensions\n        if len(flow_feat.shape) == 5:  # If [batch, features, 1, 1, 1]\n            flow_feat = flow_feat.squeeze(-1).squeeze(-1).squeeze(-1)  # [batch, features]\n            flow_feat = flow_feat.unsqueeze(1)  # [batch, 1, features]\n        elif len(flow_feat.shape) == 3:  # If [batch, features, T]\n            flow_feat = flow_feat.permute(0, 2, 1)  # [batch, T, features]\n        elif len(flow_feat.shape) == 2:  # If [batch, T*features]\n            # Assuming num_frames-1 for optical flow\n            num_flow_frames = optical_flow.size(1)\n            flow_feat = flow_feat.view(batch_size, num_flow_frames, -1)  # [batch, T-1, features_per_frame]\n    \n        # Ensure at least one temporal dimension\n        if len(visual_feat.shape) == 2:  # [batch, features]\n            visual_feat = visual_feat.unsqueeze(1)  # [batch, 1, features]\n    \n        if len(flow_feat.shape) == 2:  # [batch, features]\n            flow_feat = flow_feat.unsqueeze(1)  # [batch, 1, features]\n        \n        # Create dynamic projection layers if shapes don't match expected dimensions\n        visual_dim = visual_feat.shape[2]\n        flow_dim = flow_feat.shape[2]\n        \n        # Create dynamic projection layers only once and without log messages\n        if not self.dynamic_projections_created:\n            if visual_dim != 64:\n                # Print only once during first forward pass\n                print(f\"Creating dynamic visual projection from {visual_dim} to {self.visual_proj.out_features}\")\n                self.visual_proj = nn.Linear(visual_dim, self.visual_proj.out_features).to(visual_feat.device)\n                \n            if flow_dim != 32:\n                # Print only once during first forward pass\n                print(f\"Creating dynamic flow projection from {flow_dim} to {self.flow_proj.out_features}\")\n                self.flow_proj = nn.Linear(flow_dim, self.flow_proj.out_features).to(flow_feat.device)\n            \n            # Set flag to avoid repeated messages\n            self.dynamic_projections_created = True\n        else:\n            # On subsequent calls, just make sure the dimensions match without printing\n            if visual_dim != 64 and visual_dim != self.visual_proj.in_features:\n                self.visual_proj = nn.Linear(visual_dim, self.visual_proj.out_features).to(visual_feat.device)\n            \n            if flow_dim != 32 and flow_dim != self.flow_proj.in_features:\n                self.flow_proj = nn.Linear(flow_dim, self.flow_proj.out_features).to(flow_feat.device)\n    \n        # Project features to the same dimensionality\n        visual_proj = self.visual_proj(visual_feat)  # [batch, T, output_dim]\n        flow_proj = self.flow_proj(flow_feat)  # [batch, T, output_dim]\n    \n        # Simple combination if temporal dimensions are different\n        if visual_proj.size(1) != flow_proj.size(1):\n            # Compute temporal average for both features\n            visual_proj_avg = visual_proj.mean(dim=1, keepdim=True)  # [batch, 1, output_dim]  \n            flow_proj_avg = flow_proj.mean(dim=1, keepdim=True)  # [batch, 1, output_dim]\n        \n            # Concatenate averaged features\n            combined_feat = torch.cat([visual_proj_avg, flow_proj_avg], dim=1)  # [batch, 2, output_dim]\n        else:\n            # Combine features as a sequence\n            combined_feat = (visual_proj + flow_proj) / 2  # [batch, T, output_dim]  \n    \n        # Apply temporal attention \n        attended_feat = self.temporal_attention(combined_feat)  # [batch, T, output_dim]\n    \n        # Temporal pooling to obtain a single representation\n        fusion_feat = attended_feat.mean(dim=1)  # [batch, output_dim]\n    \n        # Fusion layer\n        fusion_feat = self.fusion_layer(fusion_feat)\n    \n        # Apply dropout regularization  \n        fusion_feat = self.dropout(fusion_feat)\n    \n        # Classification and regression\n        score = self.classifier(fusion_feat)\n        alert_pred = self.regressor(fusion_feat)\n    \n        return score, alert_pred","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.817054Z","iopub.execute_input":"2025-04-12T23:35:28.817551Z","iopub.status.idle":"2025-04-12T23:35:28.847979Z","shell.execute_reply.started":"2025-04-12T23:35:28.817522Z","shell.execute_reply":"2025-04-12T23:35:28.847416Z"},"papermill":{"duration":0.034341,"end_time":"2025-04-12T04:06:28.927300","exception":false,"start_time":"2025-04-12T04:06:28.892959","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"201590f4","cell_type":"markdown","source":"## Training Configuration and Loss Function","metadata":{"papermill":{"duration":0.004607,"end_time":"2025-04-12T04:06:28.936691","exception":false,"start_time":"2025-04-12T04:06:28.932084","status":"completed"},"tags":[]}},{"id":"642749b2","cell_type":"code","source":"# Custom loss function for time prediction\ndef custom_time_loss(pred_time, true_time):\n    \"\"\"\n    Custom loss function for time prediction that penalizes late predictions more heavily.\n    \n    This loss function combines standard MSE with an additional penalty\n    when the predicted time is later than the ground truth time.\n    \n    Args:\n        pred_time (torch.Tensor): Predicted time values\n        true_time (torch.Tensor): Ground truth time values\n        \n    Returns:\n        torch.Tensor: Combined loss value\n    \"\"\"\n    # Calculate standard MSE loss\n    base_loss = F.mse_loss(pred_time, true_time)\n    \n    # Identify and penalize late predictions (when pred_time > true_time)\n    late_mask = (pred_time > true_time).float()\n    late_penalty = F.mse_loss(pred_time * late_mask, true_time * late_mask) * 1.5\n    \n    # Combine base loss with late prediction penalty\n    return base_loss + late_penalty\n\nclass FocalLoss(nn.Module):\n    def __init__(self, alpha_init=0.25, gamma_init=2.0):\n        super(FocalLoss, self).__init__()\n        self.alpha = alpha_init\n        self.gamma = gamma_init\n        self.momentum = 0.9\n        self.batch_history = []  # Track the recent difficulty history of batches\n        \n    def forward(self, inputs, targets):\n        BCE_loss = F.binary_cross_entropy(inputs, targets, reduction='none')\n        pt = torch.exp(-BCE_loss)\n        \n        # More detailed analysis of difficulty\n        batch_difficulty = 1 - pt.mean().item()\n        self.batch_history.append(batch_difficulty)\n        \n        # Dynamically adjust gamma based on recent history\n        if len(self.batch_history) > 5:\n            recent_difficulty = sum(self.batch_history[-5:]) / 5\n            trend = (recent_difficulty - sum(self.batch_history[-10:-5]) / 5) if len(self.batch_history) > 10 else 0\n            \n            # Increase gamma for hard examples with a worsening trend\n            if recent_difficulty > 0.4 and trend > 0:\n                target_gamma = 3.0\n            # Decrease gamma for easy examples with an improving trend\n            elif recent_difficulty < 0.3 and trend < 0:\n                target_gamma = 1.5\n            else:\n                target_gamma = 2.0 + recent_difficulty\n                \n            self.gamma = self.momentum * self.gamma + (1 - self.momentum) * target_gamma\n        else:\n            self.gamma = self.momentum * self.gamma + (1 - self.momentum) * (2.0 + batch_difficulty)\n        \n        # Adjust alpha to better balance between classes\n        pos_ratio = targets.mean().item()\n        self.alpha = max(0.2, min(0.8, 1 - pos_ratio))  # Give more weight to the minority class\n        \n        F_loss = self.alpha * (1 - pt)**self.gamma * BCE_loss\n        return F_loss.mean()\n\n# Optimized training function\ndef train_model(model, dataloader, val_dataloader=None, num_epochs=30, lr=5e-5, device='cuda'):\n    \"\"\"\n    Train a deep learning model with combined classification and regression tasks,\n    including sample difficulty analysis and early stopping.\n    \n    Args:\n        model: Neural network model to train\n        dataloader: DataLoader containing training data\n        val_dataloader: Optional DataLoader for validation data\n        num_epochs: Number of training epochs\n        lr: Learning rate for optimizer\n        device: Device to run training on ('cuda' or 'cpu')\n    \n    Returns:\n        The trained model\n    \"\"\"\n    # Move model to specified device (GPU/CPU)\n    model = model.to(device)\n    \n    # Initialize Adam optimizer with weight decay for regularization\n    optimizer = optim.Adam(model.parameters(), lr=lr, weight_decay=1e-4)\n    \n    # Learning rate scheduler with cosine annealing\n    scheduler = optim.lr_scheduler.CosineAnnealingWarmRestarts(\n        optimizer, \n        T_0=10,         # Restart period\n        T_mult=1,       # Multiplication factor\n        eta_min=1e-6    # Minimum learning rate\n    )\n    \n    # Define loss functions\n    criterion_cls = FocalLoss(alpha_init=0.25, gamma_init=2.0)\n    \n    # Tracking variables for early stopping\n    best_loss = float('inf')\n    best_auc = 0.0\n    patience = 5\n    no_improve = 0\n    \n    # For tracking training progress\n    epoch_losses = []\n    val_metrics = []\n    \n    for epoch in range(num_epochs):\n        model.train()  # Set model to training mode\n        running_loss = 0.0\n        batch_count = 0\n        all_losses = []  # For storing per-sample losses\n        epoch_samples = {'correct': [], 'incorrect': []}\n        \n        for i, batch in enumerate(dataloader):\n            # Transfer batch data to device (GPU/CPU)\n            frames = batch['frames'].to(device)            # Video frames\n            optical_flow = batch['optical_flow'].to(device)  # Optical flow data\n            labels = batch['label'].to(device).unsqueeze(1)  # Classification labels\n            alert_time = batch['alert_time'].to(device).unsqueeze(1)  # Regression targets\n            video_ids = batch['video_id']\n            \n            # Zero gradients before forward pass\n            optimizer.zero_grad()\n            \n            try:\n                # Forward pass through the model\n                pred_score, pred_alert = model(frames, optical_flow)\n                \n                # Calculate classification and regression losses\n                loss_cls = criterion_cls(pred_score, labels)\n                loss_reg = custom_time_loss(pred_alert, alert_time)\n                \n                # Store individual sample losses for mining\n                sample_losses = F.binary_cross_entropy(pred_score, labels, reduction='none')\n                for j, sl in enumerate(sample_losses):\n                    all_losses.append((sl.item(), video_ids[j]))\n                \n                # Combined loss\n                loss = 0.7 * loss_cls + 0.3 * loss_reg\n                \n                # Backward pass to compute gradients\n                loss.backward()\n                \n                # Gradient clipping for training stability\n                torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=1.0)\n                \n                # Update model parameters\n                optimizer.step()\n                \n                # Record loss value\n                current_loss = loss.item()\n                running_loss += current_loss\n                batch_count += 1\n                \n                # Track correctly/incorrectly classified samples\n                binary_preds = (pred_score > 0.5).float()\n                correct_mask = (binary_preds == labels).cpu().squeeze()\n                \n                for j, (is_correct, vid) in enumerate(zip(correct_mask, video_ids)):\n                    if is_correct:\n                        epoch_samples['correct'].append(vid)\n                    else:\n                        epoch_samples['incorrect'].append(vid)\n                \n                # Free memory to prevent CUDA OOM errors\n                del frames, optical_flow, labels, alert_time, pred_score, pred_alert, loss\n                if device == 'cuda':\n                    torch.cuda.empty_cache()\n                \n                # Print progress information (reduced frequency)\n                if i % 20 == 0:\n                    print(f\"Epoch {epoch+1}/{num_epochs} - Batch {i}/{len(dataloader)} - Loss: {current_loss:.4f}\")\n                \n            except RuntimeError as e:\n                # Handle out of memory errors by skipping the batch\n                if 'out of memory' in str(e):\n                    print('| WARNING: ran out of memory, skipping batch')\n                    if device == 'cuda':\n                        torch.cuda.empty_cache()\n                else:\n                    raise e\n            except Exception as e:\n                # Handle other exceptions during training\n                print(f\"Error during training at batch {i}: {str(e)}\")\n                continue\n        \n        # At the end of each epoch, identify difficult samples\n        if epoch > 0 and epoch % 5 == 0:  # Every 5 epochs\n            difficult_samples = sorted(all_losses, key=lambda x: x[0], reverse=True)[:100]\n            print(f\"Top 10 most difficult samples: {difficult_samples[:10]}\")\n            \n            # Analyze the distribution of difficult samples\n            print(f\"Difficult samples count: {len(difficult_samples)}\")\n            difficult_ids = [ds[1] for ds in difficult_samples]\n            print(f\"Sample difficulty distribution analysis complete\")\n        \n        # Calculate average loss for the epoch\n        epoch_loss = running_loss / max(1, batch_count)\n        epoch_losses.append(epoch_loss)\n        \n        # Validation phase\n        if val_dataloader is not None:\n            metrics = validate_model(model, val_dataloader, device)\n            val_metrics.append(metrics)\n            current_auc = metrics['auc']\n            \n            print(f\"Epoch {epoch+1}/{num_epochs} - Loss: {epoch_loss:.4f}, Val AUC: {current_auc:.4f}\")\n            print(f\"Val Precision: {metrics['precision']:.4f}, Recall: {metrics['recall']:.4f}, F1: {metrics['f1']:.4f}\")\n            \n            # Early stopping based on AUC\n            if current_auc > best_auc:\n                best_auc = current_auc\n                torch.save({\n                    'epoch': epoch,\n                    'model_state_dict': model.state_dict(),\n                    'optimizer_state_dict': optimizer.state_dict(),\n                    'auc': current_auc,\n                }, \"best_model_auc.pth\")\n                print(f\"Model saved with AUC {current_auc:.4f}\")\n                no_improve = 0\n            else:\n                no_improve += 1\n                if no_improve >= patience:\n                    print(f\"Early stopping triggered at epoch {epoch+1}\")\n                    break\n        else:\n            print(f\"Epoch {epoch+1}/{num_epochs} - Loss: {epoch_loss:.4f}\")\n            \n            # Early stopping based on loss\n            if epoch_loss < best_loss:\n                best_loss = epoch_loss\n                torch.save({\n                    'epoch': epoch,\n                    'model_state_dict': model.state_dict(),\n                    'optimizer_state_dict': optimizer.state_dict(),\n                    'loss': epoch_loss,\n                }, \"best_model.pth\")\n                print(f\"Model saved with loss {epoch_loss:.4f}\")\n                no_improve = 0\n            else:\n                no_improve += 1\n                if no_improve >= patience:\n                    print(f\"Early stopping triggered at epoch {epoch+1}\")\n                    break\n        \n        # Update learning rate based on scheduler\n        scheduler.step()\n        \n        # Print training statistics\n        print(f\"Epoch {epoch+1} - Correctly classified: {len(epoch_samples['correct'])}\")\n        print(f\"Epoch {epoch+1} - Incorrectly classified: {len(epoch_samples['incorrect'])}\")\n        \n        # Save checkpoint after each epoch\n        torch.save({\n            'epoch': epoch,\n            'model_state_dict': model.state_dict(),\n            'optimizer_state_dict': optimizer.state_dict(),\n            'loss': epoch_loss,\n        }, f\"checkpoint_epoch_{epoch+1}.pth\")\n        \n        # Generate progress plot every 5 epochs\n        if (epoch + 1) % 5 == 0:\n            plot_training_progress(epoch_losses, val_metrics, f\"training_progress_epoch_{epoch+1}.png\")\n    \n    # Generate final training progress plot\n    plot_training_progress(epoch_losses, val_metrics, \"final_training_progress.png\")\n    \n    return model\n\ndef validate_model(model, val_dataloader, device='cuda'):\n    \"\"\"\n    Validates the model and returns performance metrics.\n    \n    Args:\n        model: Model to validate\n        val_dataloader: Validation data loader\n        device: Device to run validation on\n        \n    Returns:\n        dict: Dictionary containing validation metrics\n    \"\"\"\n    model.eval()\n    all_preds = []\n    all_labels = []\n    \n    with torch.no_grad():\n        for batch in val_dataloader:\n            frames = batch['frames'].to(device)\n            optical_flow = batch['optical_flow'].to(device)\n            labels = batch['label']\n            \n            scores, _ = model(frames, optical_flow)\n            \n            all_preds.extend(scores.cpu().numpy())\n            all_labels.extend(labels.numpy())\n            \n            del frames, optical_flow, scores\n            if device == 'cuda':\n                torch.cuda.empty_cache()\n    \n    # Convert to numpy arrays\n    y_pred = np.array([p[0] for p in all_preds])\n    y_true = np.array(all_labels)\n    \n    # Calculate metrics\n    auc = roc_auc_score(y_true, y_pred)\n    y_pred_binary = (y_pred >= 0.5).astype(int)\n    \n    # Compute precision, recall and F1\n    tp = np.sum((y_pred_binary == 1) & (y_true == 1))\n    fp = np.sum((y_pred_binary == 1) & (y_true == 0))\n    tn = np.sum((y_pred_binary == 0) & (y_true == 0))\n    fn = np.sum((y_pred_binary == 0) & (y_true == 1))\n    \n    precision = tp / max(tp + fp, 1)\n    recall = tp / max(tp + fn, 1)\n    f1 = 2 * (precision * recall) / max(precision + recall, 1e-10)\n    \n    return {\n        'auc': auc,\n        'precision': precision,\n        'recall': recall,\n        'f1': f1,\n        'tp': tp,\n        'fp': fp,\n        'tn': tn,\n        'fn': fn\n    }\n\ndef analyze_predictions(model, val_dataloader, device='cuda'):\n    \"\"\"\n    Analyzes model performance on the validation set.\n\n    Parameters:\n        model (nn.Module): The trained model.\n        val_dataloader (DataLoader): DataLoader for the validation set.\n        device (str): The device to run inference on ('cuda' or 'cpu').\n\n    Returns:\n        auc (float): ROC AUC score on the validation set.\n        false_positives (list): List of (video_id, prediction) for top false positives.\n        false_negatives (list): List of (video_id, prediction) for top false negatives.\n    \"\"\"\n    model.eval()\n    predictions = []\n    ground_truths = []\n    video_ids = []\n    \n    with torch.no_grad():\n        for batch in val_dataloader:\n            frames = batch['frames'].to(device)\n            optical_flow = batch['optical_flow'].to(device)\n            labels = batch['label']\n            batch_video_ids = batch['video_id']\n            \n            scores, _ = model(frames, optical_flow)\n            \n            predictions.extend(scores.cpu().numpy())\n            ground_truths.extend(labels.numpy())\n            video_ids.extend(batch_video_ids)\n            \n            del frames, optical_flow, scores\n            if device == 'cuda':\n                torch.cuda.empty_cache()\n    \n    # Convert to numpy arrays\n    predictions = np.array([p[0] for p in predictions])\n    ground_truths = np.array(ground_truths)\n    \n    # Compute ROC AUC\n    auc = roc_auc_score(ground_truths, predictions)\n    print(f\"Validation AUC: {auc:.4f}\")\n    \n    # Analyze performance at various thresholds\n    thresholds = [0.3, 0.4, 0.5, 0.6, 0.7]\n    for threshold in thresholds:\n        binary_preds = (predictions >= threshold).astype(int)\n        \n        # Compute metrics\n        tp = np.sum((binary_preds == 1) & (ground_truths == 1))\n        fp = np.sum((binary_preds == 1) & (ground_truths == 0))\n        tn = np.sum((binary_preds == 0) & (ground_truths == 0))\n        fn = np.sum((binary_preds == 0) & (ground_truths == 1))\n        \n        precision = tp / max(tp + fp, 1)\n        recall = tp / max(tp + fn, 1)\n        f1 = 2 * (precision * recall) / max(precision + recall, 1e-10)\n        \n        print(f\"Threshold {threshold:.1f} - Precision: {precision:.4f}, Recall: {recall:.4f}, F1: {f1:.4f}\")\n    \n    # Analyze errors\n    binary_preds = (predictions >= 0.5).astype(int)\n    \n    # False positives (predicted collision, but no collision occurred)\n    false_positives = [(vid, pred) for vid, pred, gt in \n                       zip(video_ids, predictions, ground_truths) \n                       if (pred >= 0.5 and gt == 0)]\n    \n    # False negatives (no collision predicted, but a collision occurred)\n    false_negatives = [(vid, pred) for vid, pred, gt in \n                       zip(video_ids, predictions, ground_truths) \n                       if (pred < 0.5 and gt == 1)]\n    \n    print(f\"\\nTop 10 most confident false positives: {sorted(false_positives, key=lambda x: x[1], reverse=True)[:10]}\")\n    print(f\"Top 10 least confident false negatives: {sorted(false_negatives, key=lambda x: x[1])[:10]}\")\n    \n    return auc, false_positives, false_negatives\n\ndef analyze_video_features(video_dir, annotations, num_videos=20):\n    \"\"\"\n    Analyzes video characteristics to better understand challenges in the dataset.\n\n    Parameters:\n        video_dir (str): Path to the directory containing video files.\n        annotations (dict): Dictionary mapping video IDs to metadata including labels.\n        num_videos (int): Number of videos to analyze (default is 20).\n\n    Returns:\n        results (dict): Dictionary containing lists of statistics for positive and negative samples.\n    \"\"\"\n    results = {'positive': [], 'negative': []}\n    \n    for video_id, data in list(annotations.items())[:num_videos]:\n        try:\n            video_path = os.path.join(video_dir, f\"{video_id}.mp4\")\n            if not os.path.exists(video_path):\n                continue\n                \n            # Extract frames for analysis\n            cap = cv2.VideoCapture(video_path)\n            if not cap.isOpened():\n                continue\n                \n            total_frames = int(cap.get(cv2.CAP_PROP_FRAME_COUNT))\n            \n            # Metrics to track\n            brightness_values = []\n            motion_values = []\n            \n            # Analyze key frames\n            prev_frame = None\n            for frame_idx in range(0, total_frames, max(1, total_frames // 10)):  # Sample 10 frames\n                cap.set(cv2.CAP_PROP_POS_FRAMES, frame_idx)\n                ret, frame = cap.read()\n                if not ret:\n                    continue\n                    \n                # Compute brightness\n                gray = cv2.cvtColor(frame, cv2.COLOR_BGR2GRAY)\n                brightness = np.mean(gray)\n                brightness_values.append(brightness)\n                \n                # Compute motion (if previous frame is available)\n                if prev_frame is not None:\n                    prev_gray = cv2.cvtColor(prev_frame, cv2.COLOR_BGR2GRAY)\n                    \n                    # Basic optical flow\n                    flow = cv2.calcOpticalFlowFarneback(prev_gray, gray, None,\n                                                        0.5, 3, 15, 3, 5, 1.2, 0)\n                    mag, _ = cv2.cartToPolar(flow[..., 0], flow[..., 1])\n                    motion = np.mean(mag)\n                    motion_values.append(motion)\n                \n                prev_frame = frame.copy()\n            \n            cap.release()\n            \n            # Compute statistics\n            avg_brightness = np.mean(brightness_values) if brightness_values else 0\n            std_brightness = np.std(brightness_values) if brightness_values else 0\n            avg_motion = np.mean(motion_values) if motion_values else 0\n            std_motion = np.std(motion_values) if motion_values else 0\n            \n            result = {\n                'video_id': video_id,\n                'avg_brightness': avg_brightness,\n                'std_brightness': std_brightness,\n                'avg_motion': avg_motion,\n                'std_motion': std_motion,\n                'total_frames': total_frames\n            }\n            \n            # Append to positive or negative results\n            if data['label'] == 1:\n                results['positive'].append(result)\n            else:\n                results['negative'].append(result)\n                \n        except Exception as e:\n            print(f\"Error analyzing video {video_id}: {str(e)}\")\n    \n    # Show comparative statistics\n    print(\"=== Positive vs. Negative Video Comparison ===\")\n    \n    pos_brightness = np.mean([r['avg_brightness'] for r in results['positive']])\n    neg_brightness = np.mean([r['avg_brightness'] for r in results['negative']])\n    print(f\"Average Brightness - Positive: {pos_brightness:.2f}, Negative: {neg_brightness:.2f}\")\n    \n    pos_motion = np.mean([r['avg_motion'] for r in results['positive']])\n    neg_motion = np.mean([r['avg_motion'] for r in results['negative']])\n    print(f\"Average Motion - Positive: {pos_motion:.2f}, Negative: {neg_motion:.2f}\")\n    \n    return results\n\ndef plot_training_progress(loss_history, val_metrics, output_path='training_progress.png'):\n    \"\"\"\n    Creates a visualization of training progress over epochs.\n\n    Parameters:\n        loss_history (list of float): List of training loss values for each epoch.\n        val_metrics (list of dict): List of validation metrics per epoch. Each dict should contain\n                                    'auc', 'precision', 'recall', and 'f1' keys.\n        output_path (str): Path where the plot image will be saved.\n\n    Saves:\n        A PNG file showing training loss and validation metrics over time.\n    \"\"\"\n    # Create the figure\n    plt.figure(figsize=(12, 8))\n    \n    # Plot training loss\n    plt.subplot(2, 1, 1)\n    plt.plot(loss_history, 'b-', label='Training Loss')\n    plt.xlabel('Epoch')\n    plt.ylabel('Loss')\n    plt.title('Training Loss')\n    plt.grid(True)\n    plt.legend()\n    \n    # Plot validation metrics\n    if val_metrics:\n        epochs = list(range(len(val_metrics)))\n        auc_values = [metrics['auc'] for metrics in val_metrics]\n        prec_values = [metrics['precision'] for metrics in val_metrics]\n        rec_values = [metrics['recall'] for metrics in val_metrics]\n        f1_values = [metrics['f1'] for metrics in val_metrics]\n        \n        plt.subplot(2, 1, 2)\n        plt.plot(epochs, auc_values, 'g-', label='AUC')\n        plt.plot(epochs, prec_values, 'r-', label='Precision')\n        plt.plot(epochs, rec_values, 'c-', label='Recall')\n        plt.plot(epochs, f1_values, 'm-', label='F1 Score')\n        plt.xlabel('Epoch')\n        plt.ylabel('Score')\n        plt.title('Validation Metrics')\n        plt.grid(True)\n        plt.legend()\n    \n    plt.tight_layout()\n    plt.savefig(output_path)\n    plt.close()\n    \n    print(f\"Progress plot saved to {output_path}\")","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.848754Z","iopub.execute_input":"2025-04-12T23:35:28.849006Z","iopub.status.idle":"2025-04-12T23:35:28.889938Z","shell.execute_reply.started":"2025-04-12T23:35:28.848978Z","shell.execute_reply":"2025-04-12T23:35:28.889362Z"},"papermill":{"duration":0.018099,"end_time":"2025-04-12T04:06:28.959289","exception":false,"start_time":"2025-04-12T04:06:28.941190","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"50de4165","cell_type":"markdown","source":"## Model Evaluation and Prediction Generation","metadata":{"papermill":{"duration":0.004374,"end_time":"2025-04-12T04:06:28.968223","exception":false,"start_time":"2025-04-12T04:06:28.963849","status":"completed"},"tags":[]}},{"id":"dd076e7c","cell_type":"code","source":"def evaluate_model(model, dataloader, device='cuda'):\n    model = model.to(device)\n    model.eval()\n    predictions = []\n    with torch.no_grad():\n        for batch in dataloader:\n            frames = batch['frames'].to(device)\n            optical_flow = batch['optical_flow'].to(device)\n            video_ids = batch.get('video_id', None)\n            pred_score, _ = model(frames, optical_flow)\n            predictions.extend(pred_score.cpu().numpy())\n    return predictions\n\ndef predict_with_dynamic_threshold(model, frames, optical_flow, base_threshold=0.5, device='cuda'):\n    \"\"\"\n    Makes predictions using a dynamic threshold adjusted to video conditions.\n    \n    Args:\n        model: Trained model\n        frames: Sequence of video frames [batch, frames, C, H, W]\n        optical_flow: Computed optical flow [batch, frames-1, 2, H, W]\n        base_threshold: Base threshold for classification\n        device: Device to run on ('cuda' or 'cpu')\n    \n    Returns:\n        prediction: Binary prediction (0 or 1)\n        score: Confidence score (0 to 1)\n        alert_time: Predicted alert time\n    \"\"\"\n    model.eval()\n    frames = frames.to(device)\n    optical_flow = optical_flow.to(device)\n    \n    # Calculate frame statistics for threshold adjustment\n    brightness = frames.float().mean().item()\n    contrast = frames.float().std().item()\n    \n    # Adjust threshold based on conditions\n    threshold = base_threshold\n    if brightness < 0.3:  # Dark video\n        threshold += 0.05  # Be more conservative in low-light conditions\n    elif contrast < 0.1:  # Low contrast video\n        threshold += 0.03  # Be more conservative in low-contrast conditions\n    \n    # Make the prediction\n    with torch.no_grad():\n        score, alert_time = model(frames, optical_flow)\n    \n    # Apply the adjusted threshold\n    prediction = (score > threshold).float()\n    \n    return prediction, score, alert_time\n\ndef predict_with_uncertainty(model, frames, optical_flow, n_samples=10, device='cuda'):\n    \"\"\"\n    Makes predictions with uncertainty quantification using Monte Carlo Dropout.\n    \n    Args:\n        model: Trained model with dropout\n        frames: Sequence of video frames\n        optical_flow: Computed optical flow\n        n_samples: Number of samples for Monte Carlo Dropout\n        device: Device for execution ('cuda' or 'cpu')\n    \n    Returns:\n        mean_score: Mean prediction score\n        mean_time: Mean predicted alert time\n        std_score: Standard deviation of scores (measure of uncertainty)\n        std_time: Standard deviation of times (measure of uncertainty)\n    \"\"\"\n    model.train()  # Activate dropout during inference\n    predictions_score = []\n    predictions_time = []\n    \n    frames = frames.to(device)\n    optical_flow = optical_flow.to(device)\n    \n    for _ in range(n_samples):\n        with torch.no_grad():\n            score, alert_time = model(frames, optical_flow)\n            predictions_score.append(score)\n            predictions_time.append(alert_time)\n    \n    # Convert lists to tensors\n    predictions_score = torch.stack(predictions_score)\n    predictions_time = torch.stack(predictions_time)\n    \n    # Calculate mean and standard deviation\n    mean_score = predictions_score.mean(dim=0)\n    std_score = predictions_score.std(dim=0)\n    mean_time = predictions_time.mean(dim=0)\n    std_time = predictions_time.std(dim=0)\n    \n    # A high std_score indicates high uncertainty in the prediction\n    return mean_score, mean_time, std_score, std_time\n\ndef predict_batch_efficient(model, test_dataloader, device='cuda', batch_size=4):\n    \"\"\"\n    Efficiently performs batch predictions, with memory cleanup.\n    \"\"\"\n    model = model.to(device)\n    model.eval()\n    \n    predictions = []\n    video_ids = []\n    \n    with torch.no_grad():\n        for batch_idx, batch in enumerate(test_dataloader):\n            if batch_idx % 10 == 0:  # Reduzido para imprimir a cada 10 lotes\n                print(f\"Processing batch {batch_idx+1}/{len(test_dataloader)}\")\n            \n            # Transfer data to the device\n            frames = batch['frames'].to(device)\n            optical_flow = batch['optical_flow'].to(device)\n            batch_video_ids = batch['video_id']\n            \n            # Make predictions\n            scores, _ = model(frames, optical_flow)\n            \n            # Collect results\n            predictions.extend(scores.cpu().numpy())\n            video_ids.extend(batch_video_ids)\n            \n            # Explicitly free memory\n            del frames, optical_flow, scores\n            if device == 'cuda':\n                torch.cuda.empty_cache()\n            \n    return predictions, video_ids\n\ndef generate_submission(predictions, video_ids, output_path='submission.csv'):\n    \"\"\"\n    Generates the submission file in the required format.\n    \"\"\"\n    with open(output_path, 'w') as f:\n        f.write(\"id,score\\n\")\n        for vid, score in zip(video_ids, predictions):\n            f.write(f\"{vid},{score[0]:.4f}\\n\")\n    print(\"Submission generated:\", output_path)","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.890651Z","iopub.execute_input":"2025-04-12T23:35:28.890828Z","iopub.status.idle":"2025-04-12T23:35:28.907865Z","shell.execute_reply.started":"2025-04-12T23:35:28.890813Z","shell.execute_reply":"2025-04-12T23:35:28.907289Z"},"papermill":{"duration":0.01677,"end_time":"2025-04-12T04:06:28.989629","exception":false,"start_time":"2025-04-12T04:06:28.972859","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"e165c37b","cell_type":"markdown","source":"## Pipeline Execution","metadata":{"papermill":{"duration":0.004493,"end_time":"2025-04-12T04:06:28.998655","exception":false,"start_time":"2025-04-12T04:06:28.994162","status":"completed"},"tags":[]}},{"id":"580de8e0","cell_type":"code","source":"def train_model_with_cross_validation(base_dir, num_folds=2, num_epochs=50, device='cuda'):\n    \"\"\"\n    Implements cross-validation for more robust training.\n    \n    Args:\n        base_dir: Base directory with the data\n        num_folds: Number of folds for cross-validation\n        num_epochs: Number of epochs per fold\n        device: Device for execution\n    \n    Returns:\n        models: List of trained models for each fold\n    \"\"\"\n    # Load annotations\n    train_df = pd.read_csv(os.path.join(base_dir, \"train.csv\"))\n    \n    # Prepare folds\n    kf = KFold(n_splits=num_folds, shuffle=True, random_state=42)\n    \n    trained_models = []\n    fold_scores = []\n    \n    for fold, (train_idx, val_idx) in enumerate(kf.split(train_df)):\n        print(f\"\\n{'='*50}\\nFold {fold+1}/{num_folds}\\n{'='*50}\")\n        \n        # Split data into train and validation\n        train_fold = train_df.iloc[train_idx]\n        val_fold = train_df.iloc[val_idx]\n        \n        # Create annotations per fold\n        train_annotations = {}\n        for _, row in train_fold.iterrows():\n            video_id = f\"{int(row['id']):05d}\"\n            train_annotations[video_id] = {\n                'label': row['target'],\n                'alert_time': row['time_of_alert'] if not pd.isna(row['time_of_alert']) else 0,\n                'event_time': row['time_of_event'] if not pd.isna(row['time_of_event']) else 0\n            }\n        \n        val_annotations = {}\n        for _, row in val_fold.iterrows():\n            video_id = f\"{int(row['id']):05d}\"\n            val_annotations[video_id] = {\n                'label': row['target'],\n                'alert_time': row['time_of_alert'] if not pd.isna(row['time_of_alert']) else 0,\n                'event_time': row['time_of_event'] if not pd.isna(row['time_of_event']) else 0\n            }\n        \n        # Create datasets and dataloaders\n        train_video_dir = os.path.join(base_dir, \"train\")\n        \n        # Preprocess training data\n        train_video_ids = list(train_annotations.keys())\n        train_processed_data = parallel_preprocess_dataset(\n            train_video_dir, \n            train_video_ids,\n            num_frames=8,\n            num_workers=4\n        )\n        \n        # Preprocess validation data \n        val_video_ids = list(val_annotations.keys())\n        val_processed_data = parallel_preprocess_dataset(\n            train_video_dir,\n            val_video_ids, \n            num_frames=8,\n            num_workers=4\n        )\n        \n        # Get transformations\n        transforms = get_video_transforms()\n        \n        # Create datasets\n        train_dataset = PreprocessedDashcamDataset(\n            train_processed_data,\n            train_annotations, \n            transform=transforms['train']\n        )\n        \n        val_dataset = PreprocessedDashcamDataset(\n            val_processed_data,\n            val_annotations,\n            transform=transforms['val'] \n        )\n        \n        # Create dataloaders\n        train_dataloader = DataLoader(\n            train_dataset,\n            batch_size=8,\n            shuffle=True, \n            num_workers=0\n        )\n        \n        val_dataloader = DataLoader(\n            val_dataset,\n            batch_size=8,\n            shuffle=False,\n            num_workers=0  \n        )\n        \n        # Create and train model\n        model = LightweightMultiStreamModel(output_dim=64)\n        model.to(device)\n        \n        # Training\n        train_model(\n            model,\n            train_dataloader,\n            num_epochs=num_epochs,\n            lr=5e-5,\n            device=device\n        )\n        \n        # Evaluate on validation set\n        model.eval()\n        val_preds = []\n        val_labels = []\n        \n        with torch.no_grad():\n            for batch in val_dataloader:\n                frames = batch['frames'].to(device)\n                optical_flow = batch['optical_flow'].to(device)\n                labels = batch['label']\n                \n                scores, _ = model(frames, optical_flow)\n                \n                val_preds.extend(scores.cpu().numpy())\n                val_labels.extend(labels.numpy())\n                \n                del frames, optical_flow, scores\n                if device == 'cuda':\n                    torch.cuda.empty_cache()\n        \n        # Calculate AUC  \n        auc = roc_auc_score(val_labels, [p[0] for p in val_preds])\n        print(f\"Fold {fold+1} - Validation AUC: {auc:.4f}\")\n        \n        # Save model for this fold \n        torch.save({\n            'fold': fold,\n            'model_state_dict': model.state_dict(),\n            'auc': auc,\n        }, f\"model_fold_{fold+1}.pth\")\n        \n        # Perform additional fine-tuning\n        fine_tuned_model = fine_tune_model(\n            model=model,  # Use the already trained model\n            train_dataloader=train_dataloader,\n            num_epochs=5,\n            device=device  \n        )\n        \n        # Save model after fine-tuning\n        torch.save({\n            'fold': fold,\n            'model_state_dict': fine_tuned_model.state_dict(),\n            'auc': auc,  # Use the previous AUC as reference\n        }, f\"model_fold_{fold+1}_finetuned.pth\")\n        \n        trained_models.append(fine_tuned_model)\n        fold_scores.append(auc)\n    \n    print(f\"\\n{'='*50}\")\n    print(f\"Cross-Validation Complete\") \n    print(f\"Average AUC across folds: {np.mean(fold_scores):.4f}\")\n    print(f\"Fold AUCs: {fold_scores}\")\n    \n    return trained_models\n\ndef fine_tune_model(model, train_dataloader, num_epochs=5, device='cuda'):\n    # Model should already be on the correct device\n    \n    # Optimizer with a very low learning rate\n    optimizer = optim.Adam(model.parameters(), lr=3e-6, weight_decay=1e-4)\n    \n    # Cyclical learning rate for fine-tuning\n    scheduler = optim.lr_scheduler.CyclicLR(\n        optimizer,\n        base_lr=1e-6,\n        max_lr=1e-5,\n        step_size_up=len(train_dataloader) // 2,\n        cycle_momentum=False\n    )\n    \n    # Adaptive loss function\n    criterion_cls = FocalLoss(alpha_init=0.25, gamma_init=2.5)\n    \n    # Fine-tuning\n    model.train()\n    best_loss = float('inf')\n    best_model_state = None\n    \n    for epoch in range(num_epochs):\n        running_loss = 0.0\n        batch_count = 0\n        \n        for i, batch in enumerate(train_dataloader):\n            frames = batch['frames'].to(device)\n            optical_flow = batch['optical_flow'].to(device)\n            labels = batch['label'].to(device).unsqueeze(1)\n            alert_time = batch['alert_time'].to(device).unsqueeze(1)\n            \n            optimizer.zero_grad()\n            \n            try:\n                # Forward pass\n                pred_score, pred_alert = model(frames, optical_flow)\n                \n                # Calculate loss\n                loss_cls = criterion_cls(pred_score, labels) \n                loss_reg = custom_time_loss(pred_alert, alert_time)\n                loss = 0.8 * loss_cls + 0.2 * loss_reg  # Greater weight on classification\n                \n                # Backward pass\n                loss.backward()\n                \n                # Gradient clipping\n                torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=0.5)\n                \n                optimizer.step()\n                scheduler.step()  # Update learning rate every batch\n                \n                # Record loss value\n                current_loss = loss.item()\n                running_loss += current_loss\n                batch_count += 1\n                \n                # Free memory\n                del frames, optical_flow, labels, alert_time, pred_score, pred_alert, loss\n                if device == 'cuda':\n                    torch.cuda.empty_cache()\n                \n                # Print progress \n                if i % 20 == 0:\n                    print(f\"Fine-tuning - Epoch {epoch+1}/{num_epochs} - Batch {i}/{len(train_dataloader)} - Loss: {current_loss:.4f}\")\n                \n            except Exception as e:\n                print(f\"Error during fine-tuning on batch {i}: {str(e)}\")\n                continue\n        \n        # Calculate average loss for the epoch\n        epoch_loss = running_loss / max(1, batch_count)\n        print(f\"Fine-tuning - Epoch {epoch+1}/{num_epochs} - Loss: {epoch_loss:.4f}\") \n        \n        # Save best model\n        if epoch_loss < best_loss:\n            best_loss = epoch_loss\n            best_model_state = model.state_dict().copy()\n    \n    # Restore best model state\n    if best_model_state is not None:\n        model.load_state_dict(best_model_state)\n    \n    return model\n\ndef ensemble_predictions(test_dataloader, model_paths, weights=None, device='cuda'):\n    \"\"\"\n    Weighted combination of multiple models.\n    \n    Args:\n        test_dataloader: DataLoader for the test set\n        model_paths: List of paths to the trained models\n        weights: Weights for each model (if None, equal weights are used)\n        device: Device for execution (cuda/cpu)\n        \n    Returns:\n        predictions: Weighted average of predictions from all models\n        video_ids: Video IDs\n    \"\"\"\n    if weights is None:\n        weights = [1.0 / len(model_paths)] * len(model_paths)\n    \n    all_predictions = []\n    video_ids = None\n    \n    for i, path in enumerate(model_paths):\n        if not os.path.exists(path):\n            print(f\"Model {path} not found, skipping...\")\n            continue\n            \n        print(f\"Generating predictions with model from {path}\")\n        \n        # Create new model instance\n        model = LightweightMultiStreamModel(output_dim=64)\n        model.to(device)\n        \n        # Load weights from checkpoint\n        checkpoint = torch.load(path, map_location=device)\n        model.load_state_dict(checkpoint['model_state_dict'])\n        model.eval()\n        \n        # Generate predictions\n        predictions = []\n        batch_ids = []\n        \n        with torch.no_grad():\n            for batch in test_dataloader:\n                frames = batch['frames'].to(device)\n                optical_flow = batch['optical_flow'].to(device)\n                batch_video_ids = batch['video_id']\n                \n                # Make prediction\n                scores, _ = model(frames, optical_flow)\n                \n                predictions.extend(scores.cpu().numpy())\n                batch_ids.extend(batch_video_ids)\n                \n                # Free memory\n                del frames, optical_flow, scores\n                if device == 'cuda':\n                    torch.cuda.empty_cache()\n        \n        # Store predictions from this model, weighted by its importance\n        all_predictions.append([p * weights[i] for p in predictions])\n        \n        # Keep IDs (assuming they are the same for all models)\n        if video_ids is None:\n            video_ids = batch_ids\n    \n    # Calculate weighted average of predictions\n    if len(all_predictions) > 0:\n        ensemble_preds = np.zeros_like(all_predictions[0])\n        for preds in all_predictions:\n            ensemble_preds += preds\n    else:\n        print(\"No models found for ensemble!\")\n        return None, None\n    \n    return ensemble_preds, video_ids\n\ndef diverse_ensemble_predictions(test_dataloader, device='cuda'):\n    \"\"\"\n    Ensemble using diverse models with confidence-weighted voting.\n    \n    Parameters:\n        test_dataloader (DataLoader): DataLoader for the test set.\n        device (str): Device to run inference on, e.g., 'cuda' or 'cpu'.\n    \n    Returns:\n        final_predictions (list): Final ensemble predictions.\n        video_ids (list): Corresponding video IDs.\n    \"\"\"\n    # Different architectures and configurations\n    models_configs = [\n        {\"model_path\": \"base_model.pth\", \"weight\": 1.0},\n        {\"model_path\": \"finetuned_model.pth\", \"weight\": 1.5},  # Higher weight for the fine-tuned model\n        {\"model_path\": \"model_fold_1.pth\", \"weight\": 0.8},\n        {\"model_path\": \"model_fold_2.pth\", \"weight\": 0.8},\n    ]\n    \n    all_predictions = []\n    all_weights = []\n    video_ids = None\n    \n    for config in models_configs:\n        if not os.path.exists(config[\"model_path\"]):\n            print(f\"Model {config['model_path']} not found, skipping...\")\n            continue\n            \n        print(f\"Generating predictions with {config['model_path']}\")\n        \n        # Create and load the model\n        model = LightweightMultiStreamModel(output_dim=64)\n        model.to(device)\n        \n        checkpoint = torch.load(config[\"model_path\"], map_location=device)\n        model.load_state_dict(checkpoint['model_state_dict'])\n        model.eval()\n        \n        # Generate predictions with confidence estimation\n        predictions = []\n        confidences = []\n        batch_ids = []\n        \n        with torch.no_grad():\n            for batch in test_dataloader:\n                frames = batch['frames'].to(device)\n                optical_flow = batch['optical_flow'].to(device)\n                batch_video_ids = batch['video_id']\n                \n                # Use Monte Carlo Dropout to estimate uncertainty\n                model.train()  # Enable dropout for Monte Carlo sampling\n                n_samples = 5\n                pred_samples = []\n                \n                for _ in range(n_samples):\n                    scores, _ = model(frames, optical_flow)\n                    pred_samples.append(scores)\n                \n                # Compute mean and standard deviation\n                pred_samples = torch.stack(pred_samples)\n                mean_preds = pred_samples.mean(dim=0)\n                std_preds = pred_samples.std(dim=0)\n                \n                # Confidence is inversely proportional to uncertainty\n                pred_confidence = 1.0 / (1.0 + std_preds)\n                \n                # Normalize confidence to [0, 1]\n                pred_confidence = pred_confidence / pred_confidence.max()\n                \n                predictions.extend(mean_preds.cpu().numpy())\n                confidences.extend(pred_confidence.cpu().numpy())\n                batch_ids.extend(batch_video_ids)\n                \n                # Free memory\n                del frames, optical_flow, pred_samples\n                if device == 'cuda':\n                    torch.cuda.empty_cache()\n        \n        # Store predictions and sample weights\n        all_predictions.append(predictions)\n        # Combine model weight with per-sample confidence\n        sample_weights = [config[\"weight\"] * conf for conf in confidences]\n        all_weights.append(sample_weights)\n        \n        if video_ids is None:\n            video_ids = batch_ids\n    \n    # Compute the weighted average of predictions\n    final_predictions = np.zeros(len(video_ids))\n    sum_weights = np.zeros(len(video_ids))\n    \n    for preds, weights in zip(all_predictions, all_weights):\n        for j, (pred, weight) in enumerate(zip(preds, weights)):\n            final_predictions[j] += pred[0] * weight[0]\n            sum_weights[j] += weight[0]\n    \n    # Normalize by total weights\n    final_predictions = final_predictions / np.maximum(sum_weights, 1e-10)\n    \n    return final_predictions, video_ids","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.909481Z","iopub.execute_input":"2025-04-12T23:35:28.909783Z","iopub.status.idle":"2025-04-12T23:35:28.939524Z","shell.execute_reply.started":"2025-04-12T23:35:28.909757Z","shell.execute_reply":"2025-04-12T23:35:28.938905Z"},"papermill":{"duration":0.168284,"end_time":"2025-04-12T04:06:29.171566","exception":false,"start_time":"2025-04-12T04:06:29.003282","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"ab1d0096","cell_type":"code","source":"if __name__ == \"__main__\":\n    # Define the directory paths\n    base_dir = \"/kaggle/input/nexar-collision-prediction\"\n    train_video_dir = os.path.join(base_dir, \"train\")\n    test_video_dir = os.path.join(base_dir, \"test\")\n    \n    print(f\"Using training directory: {train_video_dir}\")\n    print(f\"Using testing directory: {test_video_dir}\")\n    \n    # Check GPU availability\n    device = 'cuda' if torch.cuda.is_available() else 'cpu'\n    print(f\"Using device: {device}\")\n    \n    # Define the number of workers for parallel processing\n    num_workers = 4 if device == 'cuda' else 2\n    \n    # Load and prepare annotations\n    train_df = pd.read_csv(os.path.join(base_dir, \"train.csv\"))\n    test_df = pd.read_csv(os.path.join(base_dir, \"test.csv\"))\n    \n    # Create annotation dictionaries\n    train_annotations = {}\n    for _, row in train_df.iterrows():\n        video_id = f\"{int(row['id']):05d}\"\n        train_annotations[video_id] = {\n            'label': row['target'],\n            'alert_time': row['time_of_alert'] if not pd.isna(row['time_of_alert']) else 0,\n            'event_time': row['time_of_event'] if not pd.isna(row['time_of_event']) else 0\n        }\n    \n    test_annotations = {}\n    for _, row in test_df.iterrows():\n        video_id = f\"{int(row['id']):05d}\"\n        test_annotations[video_id] = {\n            'label': 0,  # Placeholder\n            'alert_time': 0  # Placeholder\n        }\n    \n    # Get transformations for data augmentation\n    transforms = get_video_transforms()\n    \n    # STEP 1: Parallel Preprocessing of Training Videos\n    print(\"Starting preprocessing of training videos...\")\n    train_video_ids = list(train_annotations.keys())\n    \n    # Preprocess training videos in parallel\n    train_processed_data = parallel_preprocess_dataset(\n        train_video_dir, \n        train_video_ids,\n        num_frames=8,\n        num_workers=num_workers\n    )\n\n    # Create dataset with preprocessed data\n    train_dataset = PreprocessedDashcamDataset(\n        train_processed_data, \n        train_annotations,\n        transform=transforms['train']\n    )\n    \n    # Create dataloader with larger batch size\n    train_dataloader = DataLoader(\n        train_dataset, \n        batch_size=16 if device == 'cuda' else 4,  # Increased batch size from 8 to 16\n        shuffle=True,\n        num_workers=0\n    )\n    \n    # STEP 2: Instantiate and Train the Model\n    print(\"Initializing model...\")\n    model = LightweightMultiStreamModel(output_dim=64)\n    \n    # Display model size\n    total_params = sum(p.numel() for p in model.parameters())\n    print(f\"Total model parameters: {total_params:,}\")\n    \n    # Initial training\n    print(\"Starting initial training...\")\n    model.to(device)\n    train_model(\n        model, \n        train_dataloader, \n        num_epochs=50,  # Full training\n        lr=1e-4, \n        device=device\n    )\n    \n    # Save the base model after initial training\n    torch.save({\n        'model_state_dict': model.state_dict(),\n        'epoch': 50,\n    }, \"base_model.pth\")\n    \n    # STEP 3: Fine-tuning\n    print(\"Starting fine-tuning...\")\n    fine_tuned_model = fine_tune_model(\n        model=model,  # Use the already trained model\n        train_dataloader=train_dataloader,\n        num_epochs=5,  # Short fine-tuning phase\n        device=device\n    )\n    \n    # Save the fine-tuned model\n    torch.save({\n        'model_state_dict': fine_tuned_model.state_dict(),\n        'epoch': 55,  # 50 initial + 5 fine-tuning\n    }, \"finetuned_model.pth\")\n    \n    # STEP 4: Preprocessing of Test Videos\n    print(\"Starting preprocessing of test videos...\")\n    test_video_ids = list(test_annotations.keys())\n    \n    # Preprocess test videos in parallel\n    test_processed_data = parallel_preprocess_dataset(\n        test_video_dir, \n        test_video_ids,\n        num_frames=8,\n        num_workers=num_workers\n    )\n\n    # Create test dataset\n    test_dataset = PreprocessedDashcamDataset(\n        test_processed_data, \n        test_annotations,\n        transform=transforms['val']\n    )\n    \n    # Create test dataloader\n    test_dataloader = DataLoader(\n        test_dataset, \n        batch_size=16 if device == 'cuda' else 4,\n        shuffle=False,\n        num_workers=0\n    )\n    \n    # STEP 5: Generate predictions with the fine-tuned model\n    print(\"Generating predictions for the test set...\")\n    model = fine_tuned_model  # Use the fine-tuned model for predictions\n    \n    predictions, video_ids = predict_batch_efficient(\n        model, \n        test_dataloader, \n        device=device, \n        batch_size=16 if device == 'cuda' else 4\n    )\n    \n    # Convert formatted IDs to integers for submission\n    original_ids = [int(vid) for vid in video_ids]\n    \n    # Generate submission file\n    submission_path = \"submission.csv\"\n    submission_df = pd.DataFrame({\n        'id': original_ids,\n        'target': [float(p[0]) for p in predictions]\n    })\n    \n    # Ensure the IDs are in the correct order\n    submission_df = submission_df.sort_values('id')\n    submission_df.to_csv(submission_path, index=False)\n    print(f\"Submission file generated: {submission_path}\")","metadata":{"execution":{"iopub.status.busy":"2025-04-12T23:35:28.940431Z","iopub.execute_input":"2025-04-12T23:35:28.940692Z","iopub.status.idle":"2025-04-13T01:54:48.541379Z","shell.execute_reply.started":"2025-04-12T23:35:28.940673Z","shell.execute_reply":"2025-04-13T01:54:48.540313Z"},"papermill":{"duration":8969.849235,"end_time":"2025-04-12T06:35:59.026207","exception":false,"start_time":"2025-04-12T04:06:29.176972","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"ed2b8ae4","cell_type":"code","source":"submission_df = pd.read_csv('/kaggle/working/submission.csv')\n \nprint(submission_df.shape)\nprint(submission_df.head())","metadata":{"execution":{"iopub.status.busy":"2025-04-13T01:54:48.542605Z","iopub.execute_input":"2025-04-13T01:54:48.542968Z","iopub.status.idle":"2025-04-13T01:54:48.552773Z","shell.execute_reply.started":"2025-04-13T01:54:48.542925Z","shell.execute_reply":"2025-04-13T01:54:48.551997Z"},"papermill":{"duration":0.025212,"end_time":"2025-04-12T06:35:59.067174","exception":false,"start_time":"2025-04-12T06:35:59.041962","status":"completed"},"tags":[],"trusted":true},"outputs":[],"execution_count":null},{"id":"80399f37-7d1b-4a7c-856c-7cce6db1a487","cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null}]}