{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.13","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[],"dockerImageVersionId":28755,"isInternetEnabled":false,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"code","source":"!pip install ultralytics","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-08-24T07:19:04.751927Z","iopub.execute_input":"2026-08-24T07:19:04.752339Z","iopub.status.idle":"2026-08-24T07:19:12.706661Z","shell.execute_reply.started":"2026-08-24T07:19:04.7523Z","shell.execute_reply":"2026-08-24T07:19:12.705658Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"import os\nimport shutil\nimport zipfile\nfrom concurrent.futures import ProcessPoolExecutor, as_completed\n\nimport cv2\nimport numpy as np\nimport pandas as pd\nimport pydicom\nfrom tqdm import tqdm\n\n\n# ============================================================\n# 1. PATHS AND CONSTANTS\n# ============================================================\n\nCOMPETITION_DATA_PATH = \"/kaggle/input/competitions/rsna-knee-abnormality-detection\"\n\nLABEL_PATH = [\n    \"/kaggle/input/notebooks/abhinav4516788/rsna-phase-1-nlp-pseudo-labels/train_pseudo_labeled_part1.csv\",\n    \"/kaggle/input/notebooks/abhinav4516788/rsna-phase1-part2/train_pseudo_labeled_part2.csv\",\n    \"/kaggle/input/notebooks/abhinav4516788/rsna-phase1-part3/train_pseudo_labeled_part3.csv\",\n    \"/kaggle/input/notebooks/abhinavreddy6480/rsna-phase1-part4/train_pseudo_labeled_part4.csv\",\n]\n\nOUTPUT_DIR = \"/kaggle/working/processed_train_arrays\"\nIMG_SIZE = 256\nNUM_SLICES = 15\nBATCH_SIZE = 1000\nNUM_WORKERS = max(1, (os.cpu_count() or 4) - 1)\n\nTARGET_CROP_MM = 150.0\nCONTEXT_SCALE = 1.15\n\nos.makedirs(OUTPUT_DIR, exist_ok=True)\n\n\n# ============================================================\n# 2. LOAD PSEUDO LABELS\n# ============================================================\n\nprint(\"Checking pseudo-label files...\\n\")\n\nfor path in LABEL_PATH:\n    if os.path.exists(path):\n        print(\"✓ Found:\", path)\n    else:\n        print(\"✗ MISSING:\", path)\n\nprint()\n\nlabels = pd.concat(\n    [pd.read_csv(path) for path in LABEL_PATH],\n    ignore_index=True\n)\n\nprint(f\"Combined label dataset: {len(labels):,} rows\")\n\nif \"StudyInstanceUID\" not in labels.columns:\n    raise ValueError(\n        \"StudyInstanceUID column was not found in the pseudo-label files.\"\n    )\n\nprint(\n    f\"Total unique labeled studies: \"\n    f\"{labels['StudyInstanceUID'].nunique():,}\"\n)\n\nprint(\"Label columns:\")\nprint(labels.columns.tolist())\nprint()\n\n\n# ============================================================\n# 3. LATERALITY NORMALIZATION\n# ============================================================\n\ndef normalize_laterality(image: np.ndarray, dcm: pydicom.Dataset) -> np.ndarray:\n    \"\"\"\n    Flip left knees horizontally so left/right anatomy\n    is presented in a consistent orientation.\n    \"\"\"\n    laterality = getattr(\n        dcm,\n        \"ImageLaterality\",\n        getattr(dcm, \"Laterality\", None)\n    )\n\n    if str(laterality).upper() == \"L\":\n        return np.fliplr(image)\n\n    return image\n\n\n# ============================================================\n# 4. KNEE ANATOMICAL LOCALIZATION\n# ============================================================\n\ndef estimate_knee_center(image: np.ndarray):\n    \"\"\"\n    Estimate the anatomical center of the knee using\n    foreground/tissue detection.\n    \"\"\"\n\n    if image is None or image.size == 0:\n        return None\n\n    h, w = image.shape[:2]\n\n    normalized = cv2.normalize(\n        image,\n        None,\n        0,\n        255,\n        cv2.NORM_MINMAX,\n        dtype=cv2.CV_8U\n    )\n\n    normalized = cv2.GaussianBlur(\n        normalized,\n        (5, 5),\n        0\n    )\n\n    threshold_value = max(\n        10,\n        int(np.percentile(normalized, 5))\n    )\n\n    _, mask = cv2.threshold(\n        normalized,\n        threshold_value,\n        255,\n        cv2.THRESH_BINARY\n    )\n\n    kernel = np.ones((5, 5), np.uint8)\n\n    mask = cv2.morphologyEx(\n        mask,\n        cv2.MORPH_OPEN,\n        kernel\n    )\n\n    mask = cv2.morphologyEx(\n        mask,\n        cv2.MORPH_CLOSE,\n        kernel\n    )\n\n    contours, _ = cv2.findContours(\n        mask,\n        cv2.RETR_EXTERNAL,\n        cv2.CHAIN_APPROX_SIMPLE\n    )\n\n    if not contours:\n        return w / 2.0, h / 2.0\n\n    min_area = 0.02 * h * w\n\n    valid_contours = [\n        contour\n        for contour in contours\n        if cv2.contourArea(contour) >= min_area\n    ]\n\n    if not valid_contours:\n        valid_contours = contours\n\n    contour = max(\n        valid_contours,\n        key=cv2.contourArea\n    )\n\n    moments = cv2.moments(contour)\n\n    if moments[\"m00\"] > 0:\n        cx = moments[\"m10\"] / moments[\"m00\"]\n        cy = moments[\"m01\"] / moments[\"m00\"]\n        return cx, cy\n\n    x, y, w_box, h_box = cv2.boundingRect(contour)\n\n    return (\n        x + w_box / 2.0,\n        y + h_box / 2.0\n    )\n\n\ndef get_series_localization_center(\n    image_paths,\n    localization_indices\n):\n    \"\"\"\n    Estimate a stable knee center for the whole MRI series\n    using several representative slices and their median center.\n    \"\"\"\n\n    centers_x = []\n    centers_y = []\n\n    for idx in localization_indices:\n        idx = max(\n            0,\n            min(len(image_paths) - 1, int(idx))\n        )\n\n        try:\n            dcm = pydicom.dcmread(\n                image_paths[idx]\n            )\n\n            image = dcm.pixel_array.astype(\n                np.float32\n            )\n\n            image = normalize_laterality(\n                image,\n                dcm\n            )\n\n            center = estimate_knee_center(\n                image\n            )\n\n            if center is not None:\n                cx, cy = center\n                centers_x.append(cx)\n                centers_y.append(cy)\n\n        except Exception:\n            continue\n\n    if not centers_x:\n        return None\n\n    return (\n        float(np.median(centers_x)),\n        float(np.median(centers_y))\n    )\n\n\n# ============================================================\n# 5. LOCALIZED ANATOMICAL CROP\n# ============================================================\n\ndef crop_localized_knee(\n    image: np.ndarray,\n    dcm: pydicom.Dataset,\n    center,\n    target_mm: float = TARGET_CROP_MM,\n    context_scale: float = CONTEXT_SCALE\n):\n    \"\"\"\n    Crop around the localized knee center.\n\n    Uses PixelSpacing when available so the physical\n    anatomical field of view is approximately consistent.\n    \"\"\"\n\n    if image is None or image.size == 0:\n        return image\n\n    h, w = image.shape[:2]\n\n    if center is None:\n        cx = w / 2.0\n        cy = h / 2.0\n    else:\n        cx, cy = center\n\n    crop_h = None\n    crop_w = None\n\n    if hasattr(dcm, \"PixelSpacing\"):\n        try:\n            spacing_y = float(dcm.PixelSpacing[0])\n            spacing_x = float(dcm.PixelSpacing[1])\n\n            if spacing_y > 0 and spacing_x > 0:\n                crop_h = int(\n                    (target_mm / spacing_y)\n                    * context_scale\n                )\n\n                crop_w = int(\n                    (target_mm / spacing_x)\n                    * context_scale\n                )\n        except Exception:\n            crop_h = None\n            crop_w = None\n\n    if crop_h is None or crop_w is None:\n        crop_h = int(h * 0.75)\n        crop_w = int(w * 0.75)\n\n    crop_size = max(\n        crop_h,\n        crop_w\n    )\n\n    crop_size = min(\n        crop_size,\n        h,\n        w\n    )\n\n    crop_size = max(\n        64,\n        crop_size\n    )\n\n    crop_size = min(\n        crop_size,\n        h,\n        w\n    )\n\n    half = crop_size // 2\n\n    x1 = int(round(cx - half))\n    y1 = int(round(cy - half))\n    x2 = x1 + crop_size\n    y2 = y1 + crop_size\n\n    if x1 < 0:\n        x2 -= x1\n        x1 = 0\n\n    if y1 < 0:\n        y2 -= y1\n        y1 = 0\n\n    if x2 > w:\n        shift = x2 - w\n        x1 -= shift\n        x2 = w\n\n    if y2 > h:\n        shift = y2 - h\n        y1 -= shift\n        y2 = h\n\n    x1 = max(0, int(x1))\n    y1 = max(0, int(y1))\n    x2 = min(w, int(x2))\n    y2 = min(h, int(y2))\n\n    cropped = image[\n        y1:y2,\n        x1:x2\n    ]\n\n    if cropped.size == 0:\n        return image\n\n    return cropped\n\n\n# ============================================================\n# 6. DICOM SORTING\n# ============================================================\n\ndef get_sorted_dicom_paths(series_folder: str) -> list:\n    \"\"\"\n    Sort DICOM files using:\n    1. ImagePositionPatient\n    2. SliceLocation\n    3. InstanceNumber\n    4. Filename\n    \"\"\"\n\n    dicom_files = [\n        os.path.join(series_folder, filename)\n        for filename in os.listdir(series_folder)\n        if filename.lower().endswith(\".dcm\")\n    ]\n\n    if not dicom_files:\n        return []\n\n    headers = []\n\n    for filepath in dicom_files:\n        try:\n            header = pydicom.dcmread(\n                filepath,\n                stop_before_pixels=True\n            )\n\n            headers.append(\n                (filepath, header)\n            )\n\n        except Exception:\n            continue\n\n    if not headers:\n        return []\n\n    if (\n        hasattr(headers[0][1], \"ImagePositionPatient\")\n        and hasattr(\n            headers[0][1],\n            \"ImageOrientationPatient\"\n        )\n    ):\n        try:\n            orientation = (\n                headers[0][1]\n                .ImageOrientationPatient\n            )\n\n            row_vector = np.array(\n                orientation[:3],\n                dtype=np.float32\n            )\n\n            col_vector = np.array(\n                orientation[3:],\n                dtype=np.float32\n            )\n\n            normal_vector = np.cross(\n                row_vector,\n                col_vector\n            )\n\n            def spatial_coordinate(item):\n                position = np.array(\n                    item[1].ImagePositionPatient,\n                    dtype=np.float32\n                )\n\n                return float(\n                    np.dot(\n                        normal_vector,\n                        position\n                    )\n                )\n\n            headers.sort(\n                key=spatial_coordinate\n            )\n\n            return [\n                filepath\n                for filepath, _ in headers\n            ]\n\n        except Exception:\n            pass\n\n    if hasattr(\n        headers[0][1],\n        \"SliceLocation\"\n    ):\n        try:\n            headers.sort(\n                key=lambda item: float(\n                    getattr(\n                        item[1],\n                        \"SliceLocation\",\n                        0\n                    )\n                )\n            )\n\n            return [\n                filepath\n                for filepath, _ in headers\n            ]\n\n        except Exception:\n            pass\n\n    if hasattr(\n        headers[0][1],\n        \"InstanceNumber\"\n    ):\n        try:\n            headers.sort(\n                key=lambda item: int(\n                    getattr(\n                        item[1],\n                        \"InstanceNumber\",\n                        0\n                    )\n                )\n            )\n\n            return [\n                filepath\n                for filepath, _ in headers\n            ]\n\n        except Exception:\n            pass\n\n    return sorted(dicom_files)\n\n\n# ============================================================\n# 7. BEST SERIES SELECTION\n# ============================================================\n\ndef get_best_series_per_study(series_df):\n    \"\"\"\n    Select the best series per study and anatomical plane.\n    \"\"\"\n\n    selected_series = []\n\n    fluid_keywords = [\n        \"fs\",\n        \"pd\",\n        \"stir\",\n        \"pdfs\"\n    ]\n\n    print(\n        \"\\nAvailable train_series.csv columns:\"\n    )\n\n    print(\n        series_df.columns.tolist()\n    )\n\n    print()\n\n    has_plane = (\n        \"Anatomical_Plane\"\n        in series_df.columns\n    )\n\n    has_description = (\n        \"SeriesDescription\"\n        in series_df.columns\n    )\n\n    print(\n        f\"Anatomical_Plane available: {has_plane}\"\n    )\n\n    print(\n        f\"SeriesDescription available: \"\n        f\"{has_description}\"\n    )\n\n    print()\n\n    if not has_plane:\n        print(\n            \"WARNING: Anatomical_Plane is missing.\"\n        )\n\n        print(\n            \"Processing all series.\"\n        )\n\n        return series_df.copy()\n\n    grouped = series_df.groupby(\n        \"StudyInstanceUID\"\n    )\n\n    for study_id, group in tqdm(\n        grouped,\n        desc=\"Selecting Best Series\"\n    ):\n        for plane in [\n            \"sagittal\",\n            \"coronal\",\n            \"axial\"\n        ]:\n            plane_mask = (\n                group[\"Anatomical_Plane\"]\n                .astype(str)\n                .str.lower()\n                .str.contains(\n                    plane[:3],\n                    na=False\n                )\n            )\n\n            plane_df = group[\n                plane_mask\n            ].copy()\n\n            if len(plane_df) == 0:\n                continue\n\n            def score_series(description):\n                if not has_description:\n                    return 0\n\n                description = str(\n                    description\n                ).lower()\n\n                score = 0\n\n                if plane in [\n                    \"sagittal\",\n                    \"coronal\"\n                ]:\n                    if any(\n                        keyword in description\n                        for keyword in fluid_keywords\n                    ):\n                        score += 10\n\n                    if \"t1\" in description:\n                        score -= 5\n\n                elif plane == \"axial\":\n                    if \"t2\" in description:\n                        score += 10\n\n                    if \"t1\" in description:\n                        score += 5\n\n                return score\n\n            if has_description:\n                plane_df[\"score\"] = (\n                    plane_df[\n                        \"SeriesDescription\"\n                    ].apply(\n                        score_series\n                    )\n                )\n            else:\n                plane_df[\"score\"] = 0\n\n            best_idx = (\n                plane_df[\"score\"]\n                .idxmax()\n            )\n\n            selected_series.append(\n                plane_df.loc[best_idx]\n            )\n\n    if not selected_series:\n        return pd.DataFrame(\n            columns=series_df.columns\n        )\n\n    return pd.DataFrame(\n        selected_series\n    ).reset_index(drop=True)\n\n\n# ============================================================\n# 8. SINGLE SERIES PROCESSOR\n# ============================================================\n\ndef process_single_series(args):\n    (\n        study_id,\n        series_id,\n        competition_path,\n        output_dir,\n        img_size\n    ) = args\n\n    save_path = os.path.join(\n        output_dir,\n        f\"{study_id}_{series_id}.npy\"\n    )\n\n    if os.path.exists(save_path):\n        return (\n            True,\n            f\"{study_id}/{series_id}\"\n        )\n\n    path_1 = os.path.join(\n        competition_path,\n        \"train_images\",\n        str(study_id),\n        str(series_id)\n    )\n\n    path_2 = os.path.join(\n        competition_path,\n        \"train_series\",\n        str(study_id),\n        str(series_id)\n    )\n\n    if os.path.exists(path_1):\n        series_folder = path_1\n    elif os.path.exists(path_2):\n        series_folder = path_2\n    else:\n        return (\n            False,\n            f\"{study_id}/{series_id} \"\n            \"(Folder not found)\"\n        )\n\n    try:\n        sorted_paths = (\n            get_sorted_dicom_paths(\n                series_folder\n            )\n        )\n\n        if not sorted_paths:\n            return (\n                False,\n                f\"{study_id}/{series_id} \"\n                \"(No valid DICOMs)\"\n            )\n\n        n = len(sorted_paths)\n\n        # ----------------------------------------------------\n        # 10% - 90% SAMPLING\n        # ----------------------------------------------------\n\n        lo = int(\n            0.1 * (n - 1)\n        )\n\n        hi = int(\n            0.9 * (n - 1)\n        )\n\n        if hi > lo:\n            indices_to_keep = (\n                np.unique(\n                    np.linspace(\n                        lo,\n                        hi,\n                        NUM_SLICES\n                    ).astype(int)\n                ).tolist()\n            )\n        else:\n            indices_to_keep = [\n                n // 2\n            ]\n\n        # ----------------------------------------------------\n        # PAD TO 15 SLICES\n        # ----------------------------------------------------\n\n        while len(indices_to_keep) < NUM_SLICES:\n            indices_to_keep.append(\n                indices_to_keep[-1]\n            )\n\n        # ----------------------------------------------------\n        # SERIES-LEVEL KNEE LOCALIZATION\n        # ----------------------------------------------------\n\n        localization_indices = (\n            np.unique(\n                np.linspace(\n                    0,\n                    len(sorted_paths) - 1,\n                    min(\n                        7,\n                        len(sorted_paths)\n                    )\n                ).astype(int)\n            ).tolist()\n        )\n\n        series_center = (\n            get_series_localization_center(\n                sorted_paths,\n                localization_indices\n            )\n        )\n\n        # ----------------------------------------------------\n        # 2.5D OVERLAPPING WINDOWS\n        # ----------------------------------------------------\n\n        resized_frames = []\n\n        for idx in indices_to_keep:\n            frame_slices = []\n\n            for offset in [-1, 0, 1]:\n                j = max(\n                    0,\n                    min(\n                        n - 1,\n                        idx + offset\n                    )\n                )\n\n                dcm = pydicom.dcmread(\n                    sorted_paths[j]\n                )\n\n                image = (\n                    dcm.pixel_array\n                    .astype(np.float32)\n                )\n\n                # --------------------------------------------\n                # LEFT / RIGHT NORMALIZATION\n                # --------------------------------------------\n\n                image = normalize_laterality(\n                    image,\n                    dcm\n                )\n\n                # --------------------------------------------\n                # ANATOMICAL LOCALIZATION\n                # --------------------------------------------\n\n                image = crop_localized_knee(\n                    image,\n                    dcm,\n                    center=series_center,\n                    target_mm=TARGET_CROP_MM,\n                    context_scale=CONTEXT_SCALE\n                )\n\n                # --------------------------------------------\n                # RESIZE\n                # --------------------------------------------\n\n                image = cv2.resize(\n                    image,\n                    (\n                        img_size,\n                        img_size\n                    ),\n                    interpolation=cv2.INTER_LINEAR\n                )\n\n                frame_slices.append(\n                    image\n                )\n\n            resized_frames.append(\n                np.stack(\n                    frame_slices,\n                    axis=0\n                )\n            )\n\n        # ----------------------------------------------------\n        # FINAL SHAPE\n        # ----------------------------------------------------\n        # (15, 3, 256, 256)\n        # ----------------------------------------------------\n\n        stacked_array = np.stack(\n            resized_frames,\n            axis=0\n        ).astype(np.float32)\n\n        # ----------------------------------------------------\n        # ROBUST INTENSITY NORMALIZATION\n        # ----------------------------------------------------\n\n        p1, p99 = np.percentile(\n            stacked_array,\n            (1, 99)\n        )\n\n        stacked_array = np.clip(\n            stacked_array,\n            p1,\n            p99\n        )\n\n        if p99 > p1:\n            stacked_array = (\n                stacked_array - p1\n            ) / (\n                p99 - p1\n            )\n        else:\n            stacked_array = (\n                np.zeros_like(\n                    stacked_array\n                )\n            )\n\n        # ----------------------------------------------------\n        # SAVE FLOAT16\n        # ----------------------------------------------------\n\n        np.save(\n            save_path,\n            stacked_array.astype(\n                np.float16\n            )\n        )\n\n        return (\n            True,\n            f\"{study_id}/{series_id}\"\n        )\n\n    except Exception as exc:\n        return (\n            False,\n            f\"{study_id}/{series_id} \"\n            f\"({str(exc)})\"\n        )\n\n\n# ============================================================\n# 9. MAIN PROCESSING\n# ============================================================\n\nif __name__ == \"__main__\":\n\n    # --------------------------------------------------------\n    # CLEAN TEMPORARY OUTPUT\n    # --------------------------------------------------------\n\n    print(\n        \"Cleaning previous temporary files...\"\n    )\n\n    shutil.rmtree(\n        OUTPUT_DIR,\n        ignore_errors=True\n    )\n\n    # --------------------------------------------------------\n    # PART 2 ZIP\n    # --------------------------------------------------------\n\n    ZIP_PATH = (\n        \"/kaggle/working/\"\n        \"rsna-15slice-part2.zip\"\n    )\n\n    if os.path.exists(ZIP_PATH):\n        os.remove(ZIP_PATH)\n\n    os.makedirs(\n        OUTPUT_DIR,\n        exist_ok=True\n    )\n\n    print(\n        \"Disk cleared. Ready to begin Part 2.\\n\"\n    )\n\n    # --------------------------------------------------------\n    # LOAD LABELS\n    # --------------------------------------------------------\n\n    df = labels.copy()\n\n    print(\n        f\"Using combined labels: \"\n        f\"{len(df):,} rows\"\n    )\n\n    # --------------------------------------------------------\n    # LOAD SERIES CSV\n    # --------------------------------------------------------\n\n    series_csv = os.path.join(\n        COMPETITION_DATA_PATH,\n        \"train_series.csv\"\n    )\n\n    series_df = pd.read_csv(\n        series_csv\n    )\n\n    print(\n        f\"Total series in competition CSV: \"\n        f\"{len(series_df):,}\"\n    )\n\n    # --------------------------------------------------------\n    # VALID STUDIES\n    # --------------------------------------------------------\n\n    valid_studies = set(\n        df[\n            \"StudyInstanceUID\"\n        ]\n        .dropna()\n        .unique()\n    )\n\n    print(\n        f\"Studies belonging to labeled data: \"\n        f\"{len(valid_studies):,}\"\n    )\n\n    # --------------------------------------------------------\n    # FILTER LABELED STUDIES\n    # --------------------------------------------------------\n\n    series_to_process = (\n        series_df[\n            series_df[\n                \"StudyInstanceUID\"\n            ].isin(valid_studies)\n        ]\n        .reset_index(drop=True)\n    )\n\n    print(\n        f\"Series belonging to labeled studies: \"\n        f\"{len(series_to_process):,}\"\n    )\n\n    # --------------------------------------------------------\n    # SELECT BEST SERIES\n    # --------------------------------------------------------\n\n    series_to_process = (\n        get_best_series_per_study(\n            series_to_process\n        )\n        .reset_index(drop=True)\n    )\n\n    print(\n        f\"\\nSeries after best-series selection: \"\n        f\"{len(series_to_process):,}\"\n    )\n\n    # ========================================================\n    # PART 2 RANGE\n    # ========================================================\n    #\n    # Part 1: 0:4500\n    # Part 2: 4500:9000\n    # Part 3: 9000:\n    #\n    # ========================================================\n\n    series_to_process = (\n        series_to_process.iloc[4500:9000]\n    )\n\n    print(\n        f\"Series selected for PART 2: \"\n        f\"{len(series_to_process):,}\"\n    )\n\n    # --------------------------------------------------------\n    # CREATE TASKS\n    # --------------------------------------------------------\n\n    tasks = [\n        (\n            row[\"StudyInstanceUID\"],\n            row[\"SeriesInstanceUID\"],\n            COMPETITION_DATA_PATH,\n            OUTPUT_DIR,\n            IMG_SIZE\n        )\n        for _, row\n        in series_to_process.iterrows()\n    ]\n\n    successful_series = 0\n    failed_series = []\n\n    print(\n        f\"\\nStarting batch preprocessing \"\n        f\"for {len(tasks):,} series \"\n        f\"using {NUM_WORKERS} workers...\"\n    )\n\n    # --------------------------------------------------------\n    # CREATE ZIP\n    # --------------------------------------------------------\n\n    with zipfile.ZipFile(\n        ZIP_PATH,\n        \"a\",\n        zipfile.ZIP_DEFLATED\n    ) as zipf:\n\n        for start_idx in range(\n            0,\n            len(tasks),\n            BATCH_SIZE\n        ):\n\n            batch_tasks = tasks[\n                start_idx:\n                start_idx + BATCH_SIZE\n            ]\n\n            batch_number = (\n                start_idx // BATCH_SIZE\n            ) + 1\n\n            print(\n                f\"\\n--- Processing Batch \"\n                f\"{batch_number} \"\n                f\"(Files {start_idx} to \"\n                f\"{min(start_idx + BATCH_SIZE, len(tasks))}) ---\"\n            )\n\n            # ------------------------------------------------\n            # MULTIPROCESSING\n            # ------------------------------------------------\n\n            with ProcessPoolExecutor(\n                max_workers=NUM_WORKERS\n            ) as executor:\n\n                futures = [\n                    executor.submit(\n                        process_single_series,\n                        task\n                    )\n                    for task in batch_tasks\n                ]\n\n                for future in tqdm(\n                    as_completed(futures),\n                    total=len(futures),\n                    desc=\"Extracting\"\n                ):\n                    try:\n                        success, identifier = (\n                            future.result()\n                        )\n\n                        if success:\n                            successful_series += 1\n                        else:\n                            failed_series.append(\n                                identifier\n                            )\n\n                    except Exception as exc:\n                        failed_series.append(\n                            f\"Worker error: {str(exc)}\"\n                        )\n\n            # ------------------------------------------------\n            # ZIP GENERATED ARRAYS\n            # ------------------------------------------------\n\n            files_to_zip = [\n                filename\n                for filename in os.listdir(\n                    OUTPUT_DIR\n                )\n                if filename.endswith(\".npy\")\n            ]\n\n            for filename in tqdm(\n                files_to_zip,\n                desc=\"Zipping & Deleting\"\n            ):\n\n                filepath = os.path.join(\n                    OUTPUT_DIR,\n                    filename\n                )\n\n                zipf.write(\n                    filepath,\n                    arcname=filename\n                )\n\n                os.remove(\n                    filepath\n                )\n\n    # ========================================================\n    # FINAL REPORT\n    # ========================================================\n\n    print(\n        \"\\n==========================================\"\n    )\n\n    print(\n        \"🎉 PART 2 FINISHED\"\n    )\n\n    print(\n        \"==========================================\"\n    )\n\n    print(\n        f\"Successfully processed: \"\n        f\"{successful_series:,}\"\n    )\n\n    print(\n        f\"Failed to process: \"\n        f\"{len(failed_series):,}\"\n    )\n\n    print(\n        f\"ZIP location: {ZIP_PATH}\"\n    )\n\n    if failed_series:\n        print(\n            \"\\nFirst 20 failed series:\"\n        )\n\n        for item in failed_series[:20]:\n            print(\n                \" -\",\n                item\n            )\n\n    # --------------------------------------------------------\n    # FINAL CLEANUP\n    # --------------------------------------------------------\n\n    shutil.rmtree(\n        OUTPUT_DIR,\n        ignore_errors=True\n    )\n\n    print(\n        \"\\nTemporary .npy files cleaned.\"\n    )\n\n    print(\n        \"ZIP file preserved.\"\n    )","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true,"execution":{"iopub.status.busy":"2026-08-24T07:19:12.709034Z","iopub.execute_input":"2026-08-24T07:19:12.709556Z"}},"outputs":[],"execution_count":null}]}