{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.6.6","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":"\"\"\"\nRSNA Intracranial Hemorrhage Detection — Preprocessing (CAPPED VERSION)\n==========================================================================\nSame as before, but processes a CAPPED, stratified subset instead of\nall 752,803 images — sized to safely fit under Kaggle's 20GB output\nquota, WITH room left over for training outputs / model checkpoints.\n\nWhy capped: the uncapped run hit the 20GB output wall at ~95%\ncompletion, causing a memory-pressure crash near the end, and only\n~448,789 images actually made it to disk before the wall was hit\n(the rest silently failed to write once the disk was full). Capping\nupfront avoids all of that — you decide the size, not the platform.\n\nSubset composition:\n- ALL positive (\"any\"=1) images — roughly ~108k, based on the\n  dataset's ~14.4% overall positive rate. Never subsample the rare\n  positive class.\n- Random negatives to fill the remainder up to TARGET_TOTAL_IMAGES.\n\nAt an observed average of ~35-45KB per processed image (varies with\nhead size / crop), 300,000 images should land around 10-14GB —\ncomfortably under the 20GB limit with meaningful headroom for\ncheckpoints (~100-300MB for ResNet-50/ViT-B/16) and CSVs.\n\nSame fixes as before, carried forward:\n- Consistent windowing (clip and rescale use the same width/level).\n- 3-channel brain/subdural/bone stack.\n- Head cropping.\n- Multiprocessing with maxtasksperchild (memory-leak fix).\n- Atomic writes with explicit format=\"PNG\".\n- Resumable — skips files that already exist.\n\"\"\"\n\nimport os\nimport shutil\nimport multiprocessing as mp\nfrom functools import partial\n\nimport numpy as np\nimport pandas as pd\nimport pydicom\nfrom PIL import Image\nfrom scipy import ndimage\nfrom tqdm import tqdm\n\n# ------------------------------------------------------------------\n# CONFIG\n# ------------------------------------------------------------------\nDATA_ROOT = \"/kaggle/input/competitions/rsna-intracranial-hemorrhage-detection/rsna-intracranial-hemorrhage-detection\"\nTRAIN_CSV = os.path.join(DATA_ROOT, \"stage_2_train.csv\")\nTRAIN_IMG_DIR = os.path.join(DATA_ROOT, \"stage_2_train\")\nOUTPUT_DIR = \"/kaggle/working/png/train_brain_subdural_bone\"\n\nNUM_WORKERS = max(mp.cpu_count() - 1, 1)\nTARGET_TOTAL_IMAGES = 300_000\nRANDOM_SEED = 42\n\nWINDOWS = {\n    \"brain\": (80, 40),\n    \"subdural\": (200, 80),\n    \"bone\": (2000, 600),\n}\n\n\ndef linear_windowing(img_hu, width, level):\n    lower = level - width / 2\n    upper = level + width / 2\n    img = np.clip(img_hu, lower, upper)\n    img = (img - lower) / (upper - lower)\n    return (img * 255).astype(np.uint8)\n\n\ndef dicom_to_hu(dcm):\n    pixels = dcm.pixel_array.astype(np.float32)\n    pixels = pixels * dcm.RescaleSlope + dcm.RescaleIntercept\n    return pixels\n\n\nclass CropHead:\n    def __init__(self, offset=10):\n        self.offset = offset\n\n    def crop_extents(self, img_bool):\n        try:\n            labeled, n = ndimage.label(img_bool)\n            if n == 0:\n                return 0, img_bool.shape[0], 0, img_bool.shape[1]\n            sizes = np.bincount(labeled.flatten())\n            sizes[0] = 0\n            head_label = np.argmax(sizes)\n            head_mask = labeled == head_label\n            rows = np.flatnonzero(head_mask.any(axis=1))\n            cols = np.flatnonzero(head_mask.any(axis=0))\n            x_min = max(rows.min() - self.offset, 0)\n            x_max = min(rows.max() + self.offset + 1, img_bool.shape[0])\n            y_min = max(cols.min() - self.offset, 0)\n            y_max = min(cols.max() + self.offset + 1, img_bool.shape[1])\n            return x_min, x_max, y_min, y_max\n        except ValueError:\n            return 0, img_bool.shape[0], 0, img_bool.shape[1]\n\n    def __call__(self, img):\n        # Collapse to 2D first — img is H x W x 3 (brain/subdural/bone\n        # channels). Finding connected components on the raw 3D stack\n        # was labeling across channels too, not just spatially, which\n        # produced inconsistent/off-center crops since each window\n        # shows different amounts of anatomy. any(axis=-1) gives a\n        # single 2D \"is there anything non-zero in ANY channel here\"\n        # mask, which is what we actually want for locating the head.\n        mask_2d = img.any(axis=-1) if img.ndim == 3 else (img > 0)\n        x_min, x_max, y_min, y_max = self.crop_extents(mask_2d)\n        return img[x_min:x_max, y_min:y_max]\n\n\ndef process_one_dicom(path, crop_head):\n    try:\n        dcm = pydicom.dcmread(path)\n    except Exception as e:\n        # Tag read failures distinctly from windowing/crop failures below,\n        # so failed-file logs make it obvious whether a file is corrupt\n        # vs. a bug in our own processing logic.\n        raise RuntimeError(f\"DICOM_READ_ERROR: {e}\") from e\n\n    hu = dicom_to_hu(dcm)\n    channels = []\n    for name in (\"brain\", \"subdural\", \"bone\"):\n        width, level = WINDOWS[name]\n        channels.append(linear_windowing(hu, width, level))\n    img = np.stack(channels, axis=-1)\n    img = crop_head(img)\n    if img.shape[0] == 0 or img.shape[1] == 0:\n        img = np.zeros((512, 512, 3), dtype=np.uint8)\n    return img\n\n\ndef _worker(fname, img_dir, output_dir):\n    out_name = fname.replace(\".dcm\", \".png\")\n    out_path = os.path.join(output_dir, out_name)\n    if os.path.exists(out_path):\n        return \"skipped\"\n    tmp_path = out_path + f\".tmp{os.getpid()}\"\n    try:\n        crop_head = CropHead()\n        img = process_one_dicom(os.path.join(img_dir, fname), crop_head)\n        Image.fromarray(img).save(tmp_path, format=\"PNG\")\n        os.replace(tmp_path, out_path)\n        return None\n    except Exception as e:\n        if os.path.exists(tmp_path):\n            try:\n                os.remove(tmp_path)\n            except OSError:\n                pass\n        return (fname, str(e))\n\n\ndef select_capped_subset(train_csv, target_total, seed=RANDOM_SEED):\n    \"\"\"\n    Returns (positive_ids, negative_ids) as two separate sets, so they\n    can be processed as separate batches — resets the multiprocessing\n    pool between batches (extra memory-leak safety on top of\n    maxtasksperchild), and keeps each zip/commit smaller and more\n    reliable than one giant 300k-file commit.\n    \"\"\"\n    y = pd.read_csv(train_csv)\n    id_split = y.ID.str.rsplit(\"_\", n=1, expand=True)\n    y = pd.concat([id_split, y.Label], axis=1)\n    y.columns = [\"id\", \"sub_type\", \"label\"]\n    y = y.drop_duplicates(subset=[\"id\", \"sub_type\"])\n\n    df = y.pivot(index=\"id\", columns=\"sub_type\", values=\"label\")\n\n    positives = df[df[\"any\"] == 1]\n    negatives = df[df[\"any\"] == 0]\n\n    n_pos = len(positives)\n    n_neg_needed = max(target_total - n_pos, 0)\n\n    print(f\"Full dataset: {len(df)} images, {n_pos} positive ({n_pos/len(df):.1%})\")\n\n    if n_neg_needed > len(negatives):\n        print(f\"WARNING: requested {n_neg_needed} negatives but only \"\n              f\"{len(negatives)} exist. Using all negatives.\")\n        n_neg_needed = len(negatives)\n\n    sampled_negatives = negatives.sample(n=n_neg_needed, random_state=seed)\n\n    print(f\"Positive batch: {n_pos} images\")\n    print(f\"Negative batch: {len(sampled_negatives)} images\")\n\n    return set(positives.index), set(sampled_negatives.index)\n\n\ndef process_id_batch(batch_ids, batch_name, all_filenames, img_dir, output_dir):\n    \"\"\"\n    Processes one batch (e.g. just positives, or just negatives) using\n    its own fresh multiprocessing pool, then zips the result.\n\n    Resumable: if this batch is interrupted and re-run, any file whose\n    output PNG already exists is skipped (see _worker) — no need to\n    restart a batch from zero after a crash.\n    \"\"\"\n    filenames = [f for f in all_filenames if f.replace(\".dcm\", \"\") in batch_ids]\n    os.makedirs(output_dir, exist_ok=True)\n    print(f\"\\n=== Batch '{batch_name}': {len(filenames)} files ===\")\n    print(f\"Estimated size at ~40KB/image: ~{len(filenames) * 40 / 1e6:.1f} GB\")\n\n    worker_fn = partial(_worker, img_dir=img_dir, output_dir=output_dir)\n    failed = []\n    skipped = 0\n    processed = 0\n\n    with mp.Pool(NUM_WORKERS, maxtasksperchild=2000) as pool:\n        for result in tqdm(pool.imap_unordered(worker_fn, filenames, chunksize=64),\n                            total=len(filenames), desc=batch_name):\n            if result == \"skipped\":\n                skipped += 1\n            elif result is None:\n                processed += 1\n            else:\n                failed.append(result)\n\n    print(f\"Batch '{batch_name}' — processed: {processed}, skipped: {skipped}, \"\n          f\"failed: {len(failed)}\")\n    if failed:\n        for f, e in failed[:5]:\n            print(f\"  {f}: {e}\")\n\n    return processed, skipped, len(failed)\n\n\nimport zipfile\n\ndef zip_output_chunked(output_dir, zip_prefix, chunk_size=25000, delete_after_each=True):\n    \"\"\"\n    Zips output in chunks instead of one giant archive — keeps peak\n    disk usage much lower, since only one chunk's worth of source\n    files + that chunk's zip need to coexist at a time, instead of\n    the ENTIRE folder + a full zip simultaneously (which is what blew\n    past the disk limit on the positive batch).\n\n    Deletes each chunk's source files right after that chunk's zip is\n    verified, so space frees up progressively as you go rather than\n    needing the whole folder to survive until the very end.\n    \"\"\"\n    files = sorted(f for f in os.listdir(output_dir) if f.endswith(\".png\"))\n    print(f\"Chunking {len(files)} files into groups of {chunk_size}...\")\n\n    for i in range(0, len(files), chunk_size):\n        chunk_files = files[i:i + chunk_size]\n        zip_path = f\"{zip_prefix}_{i // chunk_size + 1}.zip\"\n\n        with zipfile.ZipFile(zip_path, \"w\", compression=zipfile.ZIP_DEFLATED) as z:\n            for f in chunk_files:\n                z.write(os.path.join(output_dir, f), arcname=f)\n\n        if not zipfile.is_zipfile(zip_path):\n            raise RuntimeError(\n                f\"{zip_path} failed validation — NOT deleting its source \"\n                f\"files. Investigate before continuing.\"\n            )\n\n        size_gb = os.path.getsize(zip_path) / 1e9\n        print(f\"{zip_path} done ({len(chunk_files)} files, {size_gb:.2f} GB)\")\n\n        if delete_after_each:\n            for f in chunk_files:\n                os.remove(os.path.join(output_dir, f))\n            print(f\"  -> deleted {len(chunk_files)} source files, space freed\")\n\n\ndef zip_output(output_dir, zip_path, delete_source_after=True):\n    \"\"\"\n    Zips the output folder to a single file, then verifies the zip is\n    valid and deletes the source folder — otherwise folder + zip\n    coexist simultaneously, which can nearly double disk usage and\n    risk hitting the 20GB wall again (a real catch from code review:\n    for the negative batch alone, that's folder + zip both >7GB at\n    once if not cleaned up).\n\n    NOTE: for large batches, prefer zip_output_chunked() instead — this\n    single-zip version needs the ENTIRE folder plus the ENTIRE zip to\n    coexist at peak, which is what caused the positive batch's disk\n    overflow (12.8GB folder + up to ~12.8GB zip = ~25GB, past the 20GB\n    quota).\n    \"\"\"\n    print(f\"\\nZipping {output_dir} -> {zip_path} ...\")\n\n    # Quick safety check: is there enough free space to hold the zip\n    # temporarily alongside the source folder, before we've deleted it?\n    folder_size = sum(os.path.getsize(os.path.join(dp, f))\n                       for dp, _, fnames in os.walk(output_dir) for f in fnames)\n    free_space = shutil.disk_usage(output_dir).free\n    SAFETY_MARGIN = 1.2  # buffer for filesystem metadata, temp buffers,\n                          # and the fact that PNGs (already compressed)\n                          # won't shrink much further when zipped — don't\n                          # assume the zip will be meaningfully smaller\n                          # than the source folder\n    if folder_size * SAFETY_MARGIN > free_space:\n        raise RuntimeError(\n            f\"Not enough free disk space to safely zip this batch \"\n            f\"(folder is {folder_size/1e9:.2f}GB, needs \"\n            f\"{folder_size*SAFETY_MARGIN/1e9:.2f}GB with margin, only \"\n            f\"{free_space/1e9:.2f}GB free). Stopping before risking \"\n            f\"another disk-full crash.\"\n        )\n\n    base_path = zip_path.replace(\".zip\", \"\")\n    shutil.make_archive(base_path, \"zip\", output_dir)\n    size_gb = os.path.getsize(zip_path) / 1e9\n    print(f\"Zip created: {zip_path} ({size_gb:.2f} GB)\")\n\n    # Verify before deleting anything\n    if not zipfile.is_zipfile(zip_path):\n        raise RuntimeError(\n            f\"Zip at {zip_path} failed validation — NOT deleting source \"\n            f\"folder. Investigate before retrying.\"\n        )\n\n    if delete_source_after:\n        print(f\"Zip verified valid. Deleting source folder {output_dir} \"\n              f\"to free disk space before the next batch...\")\n        shutil.rmtree(output_dir)\n        os.makedirs(output_dir, exist_ok=True)  # leave it empty, ready for next batch\n        print(\"Source folder cleared.\")\n\n\ndef run_small_test(positive_ids, negative_ids, all_filenames, n=100):\n    \"\"\"\n    Processes a SMALL sample (n positives + n negatives) first, so you\n    can visually confirm the crop fix worked before committing to the\n    full batches. Look at a few output images afterward.\n    \"\"\"\n    test_pos = set(list(positive_ids)[:n])\n    test_neg = set(list(negative_ids)[:n])\n    test_output_dir = \"/kaggle/working/png/test_sample\"\n    os.makedirs(test_output_dir, exist_ok=True)\n\n    process_id_batch(test_pos | test_neg, \"small_test\",\n                      all_filenames, TRAIN_IMG_DIR, test_output_dir)\n\n    print(f\"\\nCheck a few images in {test_output_dir} now — \"\n          f\"confirm heads are NOT cut off before running the full batches.\")\n\n\ndef main_test_only():\n    \"\"\"Step 1: run this first. Small visual sanity check.\"\"\"\n    print(\"Selecting subset...\")\n    positive_ids, negative_ids = select_capped_subset(TRAIN_CSV, TARGET_TOTAL_IMAGES)\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    run_small_test(positive_ids, negative_ids, all_filenames, n=100)\n\n\ndef main_positive_batch():\n    \"\"\"Step 2: run after confirming the test sample looks correct.\"\"\"\n    positive_ids, _ = select_capped_subset(TRAIN_CSV, TARGET_TOTAL_IMAGES)\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(positive_ids, \"positive\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n    zip_output(OUTPUT_DIR, \"/kaggle/working/positive_batch.zip\")\n\n\ndef _get_sorted_positive_parts():\n    \"\"\"\n    Helper for main_positive_part2() / main_positive_part3() below.\n\n    Sorts the positive IDs deterministically and slices them into the\n    same three parts every time, so the same IDs always land in the\n    same part across runs/sessions:\n      - part1: first 25,000  (already processed and archived as\n        positive_part_1.zip — never reprocessed)\n      - part2: next 41,500\n      - part3: remainder (~41,433)\n    \"\"\"\n    positive_ids, _ = select_capped_subset(TRAIN_CSV, TARGET_TOTAL_IMAGES)\n    positive_ids = sorted(list(positive_ids))\n\n    part1 = positive_ids[:25000]\n    part2 = positive_ids[25000:66500]\n    part3 = positive_ids[66500:]\n\n    return part1, part2, part3\n\n\ndef main_positive_part2():\n    \"\"\"\n    Step 2b (alternative to main_positive_batch, used to resume after\n    Part 1 was already processed and archived as positive_part_1.zip):\n    processes only Part 2 (~41,500 positive IDs) into OUTPUT_DIR.\n\n    Does NOT zip automatically — Part 2's folder could land around\n    4.5-5GB, and given the disk-full crash last time on a bigger\n    single-zip attempt, it's safer to chunk the zip manually afterward\n    with zip_output_chunked(), e.g.:\n\n        zip_output_chunked(OUTPUT_DIR, \"/kaggle/working/positive_part_2\",\n                            chunk_size=20750)\n\n    which produces positive_part_2_1.zip, positive_part_2_2.zip, etc.\n    (rename to positive_part_2a.zip / positive_part_2b.zip afterward if\n    you want that naming). This way, if packaging fails partway through,\n    the PNGs are still on disk and preprocessing doesn't need to rerun.\n    \"\"\"\n    _, part2, _ = _get_sorted_positive_parts()\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(set(part2), \"positive_part2\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n\n\ndef main_positive_part3():\n    \"\"\"\n    Step 2c: run after Part 2 is done, zipped, and its PNGs cleared\n    from OUTPUT_DIR. Processes the remaining positive IDs (Part 3,\n    ~41,433 images) into OUTPUT_DIR.\n\n    Does NOT zip automatically, same reasoning as main_positive_part2()\n    — chunk it manually afterward with zip_output_chunked(), e.g.:\n\n        zip_output_chunked(OUTPUT_DIR, \"/kaggle/working/positive_part_3\",\n                            chunk_size=20717)\n    \"\"\"\n    _, _, part3 = _get_sorted_positive_parts()\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(set(part3), \"positive_part3\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n\n\ndef main_negative_batch():\n    \"\"\"Step 3: run after the positive batch is done and zipped.\"\"\"\n    _, negative_ids = select_capped_subset(TRAIN_CSV, TARGET_TOTAL_IMAGES)\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(negative_ids, \"negative\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n    zip_output(OUTPUT_DIR, \"/kaggle/working/negative_batch.zip\")\n\n\ndef _get_sorted_negative_parts():\n    \"\"\"\n    Helper for main_negative_part1() / main_negative_part2() /\n    main_negative_part3() below.\n\n    The sampled negative batch (192,067 images) is too large to\n    preprocess in a single run — the resulting PNGs require more disk\n    space than Kaggle provides. Sorts the sampled negative IDs\n    deterministically and slices them into three parts every time, so\n    the same IDs always land in the same part across runs/sessions:\n      - part1: first 64,000\n      - part2: next 64,000\n      - part3: remainder (~64,067)\n    \"\"\"\n    _, negative_ids = select_capped_subset(TRAIN_CSV, TARGET_TOTAL_IMAGES)\n    negative_ids = sorted(list(negative_ids))\n\n    part1 = negative_ids[:64000]\n    part2 = negative_ids[64000:128000]\n    part3 = negative_ids[128000:]\n\n    return part1, part2, part3\n\n\ndef main_negative_part1():\n    \"\"\"\n    Step 3a (alternative to main_negative_batch, used to split the\n    sampled negatives into three deterministic preprocessing runs):\n    processes only Part 1 (~64,000 negative IDs) into OUTPUT_DIR.\n\n    Does NOT zip automatically — zip manually afterward once disk\n    usage has been checked, e.g. with zip_output_chunked().\n    \"\"\"\n    part1, _, _ = _get_sorted_negative_parts()\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(set(part1), \"negative_part1\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n\n\ndef main_negative_part2():\n    \"\"\"\n    Step 3b: run after Part 1 is done, zipped, and its PNGs cleared\n    from OUTPUT_DIR. Processes Part 2 (~64,000 negative IDs) into\n    OUTPUT_DIR.\n\n    Does NOT zip automatically — zip manually afterward once disk\n    usage has been checked, e.g. with zip_output_chunked().\n    \"\"\"\n    _, part2, _ = _get_sorted_negative_parts()\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(set(part2), \"negative_part2\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n\n\ndef main_negative_part3():\n    \"\"\"\n    Step 3c: run after Part 2 is done, zipped, and its PNGs cleared\n    from OUTPUT_DIR. Processes the remaining negative IDs (Part 3,\n    ~64,067 images) into OUTPUT_DIR.\n\n    Does NOT zip automatically — zip manually afterward once disk\n    usage has been checked, e.g. with zip_output_chunked().\n    \"\"\"\n    _, _, part3 = _get_sorted_negative_parts()\n    all_filenames = os.listdir(TRAIN_IMG_DIR)\n    process_id_batch(set(part3), \"negative_part3\", all_filenames, TRAIN_IMG_DIR, OUTPUT_DIR)\n\n\nif __name__ == \"__main__\":\n    main_negative_part1()  # Nothing auto-runs. Call the functions above one at a time,\n          # each in its own cell:\n          #   1. main_test_only()          <- run first, then LOOK at\n          #                                    /kaggle/working/png/test_sample\n          #   2. main_positive_batch()      <- only after test looks good\n          #      (OR, if Part 1 is already archived, use\n          #       main_positive_part2() then main_positive_part3()\n          #       instead to resume without reprocessing Part 1 — each\n          #       one only preprocesses; call\n          #       zip_output_chunked(OUTPUT_DIR, \"/kaggle/working/positive_partN\",\n          #                          chunk_size=~20750)\n          #       manually afterward to chunk-zip and clear PNGs before\n          #       moving to the next part)\n          #   3. main_negative_batch()      <- only after all positive\n          #                                    parts are done and zipped\n          #      (OR, since the sampled negative batch is 192,067\n          #       images, use main_negative_part1(), main_negative_part2(),\n          #       then main_negative_part3() instead to split it into\n          #       three deterministic runs — each one only preprocesses;\n          #       call zip_output_chunked(OUTPUT_DIR, \"/kaggle/working/negative_partN\",\n          #                               chunk_size=~20000-ish) manually\n          #       afterward to chunk-zip and clear PNGs before moving to\n          #       the next part)","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","trusted":true},"outputs":[],"execution_count":null}]}