{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"pygments_lexer":"ipython3","nbconvert_exporter":"python","version":"3.6.4","file_extension":".py","codemirror_mode":{"name":"ipython","version":3},"name":"python","mimetype":"text/x-python"}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# RSNA Breast Baseline - Faster Inference with NVIDIA Dali\n\nTheo V to ROI pipeline, then ROI pipeline","metadata":{}},{"cell_type":"code","source":"DEBUG = False","metadata":{"execution":{"iopub.status.busy":"2023-01-13T09:48:25.985652Z","iopub.execute_input":"2023-01-13T09:48:25.986457Z","iopub.status.idle":"2023-01-13T09:48:26.009546Z","shell.execute_reply.started":"2023-01-13T09:48:25.986367Z","shell.execute_reply":"2023-01-13T09:48:26.008667Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Initialization\n- Install Dali + Overwrite a file to handle UINT16\n- Install dicomsdl, pylibjpeg","metadata":{}},{"cell_type":"code","source":"!pip install /kaggle/input/rsna-2022-whl/{pydicom-2.3.0-py3-none-any.whl,pylibjpeg-1.4.0-py3-none-any.whl,python_gdcm-3.0.15-cp37-cp37m-manylinux_2_17_x86_64.manylinux2014_x86_64.whl}\n!pip install /kaggle/input/nvidia-dali-wheel/nvidia_dali_nightly_cuda110-1.22.0.dev20221213-6757685-py3-none-manylinux2014_x86_64.whl\n!pip install /kaggle/input/nvidia-dali-wheel/dicomsdl-0.109.1-cp37-cp37m-manylinux_2_12_x86_64.manylinux2010_x86_64.whl","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19","_kg_hide-input":true,"_kg_hide-output":true,"execution":{"iopub.status.busy":"2023-01-13T09:48:26.011379Z","iopub.execute_input":"2023-01-13T09:48:26.011815Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%writefile /opt/conda/lib/python3.7/site-packages/nvidia/dali/plugin/pytorch.py\n\n# Copyright (c) 2017-2022, NVIDIA CORPORATION & AFFILIATES. All rights reserved.\n#\n# Licensed under the Apache License, Version 2.0 (the \"License\");\n# you may not use this file except in compliance with the License.\n# You may obtain a copy of the License at\n#\n#     http://www.apache.org/licenses/LICENSE-2.0\n#\n# Unless required by applicable law or agreed to in writing, software\n# distributed under the License is distributed on an \"AS IS\" BASIS,\n# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\n# See the License for the specific language governing permissions and\n# limitations under the License.\n\nfrom nvidia.dali.backend import TensorGPU, TensorListGPU\nfrom nvidia.dali.pipeline import Pipeline\nimport nvidia.dali.ops as ops\nfrom nvidia.dali import types\nfrom nvidia.dali.plugin.base_iterator import _DaliBaseIterator\nfrom nvidia.dali.plugin.base_iterator import LastBatchPolicy\nimport torch\nimport torch.utils.dlpack as torch_dlpack\nimport ctypes\nimport numpy as np\n\nto_torch_type = {\n    types.DALIDataType.FLOAT:   torch.float32,\n    types.DALIDataType.FLOAT64: torch.float64,\n    types.DALIDataType.FLOAT16: torch.float16,\n    types.DALIDataType.UINT8:   torch.uint8,\n    types.DALIDataType.INT8:    torch.int8,\n    types.DALIDataType.UINT16:  torch.int16,\n    types.DALIDataType.INT16:   torch.int16,\n    types.DALIDataType.INT32:   torch.int32,\n    types.DALIDataType.INT64:   torch.int64\n}\n\n\ndef feed_ndarray(dali_tensor, arr, cuda_stream=None):\n    \"\"\"\n    Copy contents of DALI tensor to PyTorch's Tensor.\n\n    Parameters\n    ----------\n    `dali_tensor` : nvidia.dali.backend.TensorCPU or nvidia.dali.backend.TensorGPU\n                    Tensor from which to copy\n    `arr` : torch.Tensor\n            Destination of the copy\n    `cuda_stream` : torch.cuda.Stream, cudaStream_t or any value that can be cast to cudaStream_t.\n                    CUDA stream to be used for the copy\n                    (if not provided, an internal user stream will be selected)\n                    In most cases, using pytorch's current stream is expected (for example,\n                    if we are copying to a tensor allocated with torch.zeros(...))\n    \"\"\"\n    dali_type = to_torch_type[dali_tensor.dtype]\n\n    assert dali_type == arr.dtype, (\"The element type of DALI Tensor/TensorList\"\n                                    \" doesn't match the element type of the target PyTorch Tensor: \"\n                                    \"{} vs {}\".format(dali_type, arr.dtype))\n    assert dali_tensor.shape() == list(arr.size()), \\\n        (\"Shapes do not match: DALI tensor has size {0}, but PyTorch Tensor has size {1}\".\n            format(dali_tensor.shape(), list(arr.size())))\n    cuda_stream = types._raw_cuda_stream(cuda_stream)\n\n    # turn raw int to a c void pointer\n    c_type_pointer = ctypes.c_void_p(arr.data_ptr())\n    if isinstance(dali_tensor, (TensorGPU, TensorListGPU)):\n        stream = None if cuda_stream is None else ctypes.c_void_p(cuda_stream)\n        dali_tensor.copy_to_external(c_type_pointer, stream, non_blocking=True)\n    else:\n        dali_tensor.copy_to_external(c_type_pointer)\n    return arr\n\n\nclass DALIGenericIterator(_DaliBaseIterator):\n    \"\"\"\n    General DALI iterator for PyTorch. It can return any number of\n    outputs from the DALI pipeline in the form of PyTorch's Tensors.\n\n    Parameters\n    ----------\n    pipelines : list of nvidia.dali.Pipeline\n                List of pipelines to use\n    output_map : list of str\n                List of strings which maps consecutive outputs\n                of DALI pipelines to user specified name.\n                Outputs will be returned from iterator as dictionary\n                of those names.\n                Each name should be distinct\n    size : int, default = -1\n                Number of samples in the shard for the wrapped pipeline (if there is more than\n                one it is a sum)\n                Providing -1 means that the iterator will work until StopIteration is raised\n                from the inside of iter_setup(). The options `last_batch_policy` and\n                `last_batch_padded` don't work in such case. It works with only one pipeline inside\n                the iterator.\n                Mutually exclusive with `reader_name` argument\n    reader_name : str, default = None\n                Name of the reader which will be queried to the shard size, number of shards and\n                all other properties necessary to count properly the number of relevant and padded\n                samples that iterator needs to deal with. It automatically sets `last_batch_policy`\n                to PARTIAL when the FILL is used, and `last_batch_padded` accordingly to match\n                the reader's configuration\n    auto_reset : string or bool, optional, default = False\n                Whether the iterator resets itself for the next epoch or it requires reset() to be\n                called explicitly.\n\n                It can be one of the following values:\n\n                * ``\"no\"``, ``False`` or ``None`` - at the end of epoch StopIteration is raised\n                  and reset() needs to be called\n                * ``\"yes\"`` or ``\"True\"``- at the end of epoch StopIteration is raised but reset()\n                  is called internally automatically\n\n    dynamic_shape : any, optional,\n                Parameter used only for backward compatibility.\n    fill_last_batch : bool, optional, default = None\n                **Deprecated** Please use ``last_batch_policy`` instead\n\n                Whether to fill the last batch with data up to 'self.batch_size'.\n                The iterator would return the first integer multiple\n                of self._num_gpus * self.batch_size entries which exceeds 'size'.\n                Setting this flag to False will cause the iterator to return\n                exactly 'size' entries.\n    last_batch_policy: optional, default = LastBatchPolicy.FILL\n                What to do with the last batch when there are not enough samples in the epoch\n                to fully fill it. See :meth:`nvidia.dali.plugin.base_iterator.LastBatchPolicy`\n    last_batch_padded : bool, optional, default = False\n                Whether the last batch provided by DALI is padded with the last sample\n                or it just wraps up. In the conjunction with ``last_batch_policy`` it tells\n                if the iterator returning last batch with data only partially filled with\n                data from the current epoch is dropping padding samples or samples from\n                the next epoch. If set to ``False`` next\n                epoch will end sooner as data from it was consumed but dropped. If set to\n                True next epoch would be the same length as the first one. For this to happen,\n                the option `pad_last_batch` in the reader needs to be set to True as well.\n                It is overwritten when `reader_name` argument is provided\n    prepare_first_batch : bool, optional, default = True\n                Whether DALI should buffer the first batch right after the creation of the iterator,\n                so one batch is already prepared when the iterator is prompted for the data\n\n    Example\n    -------\n    With the data set ``[1,2,3,4,5,6,7]`` and the batch size 2:\n\n    last_batch_policy = LastBatchPolicy.PARTIAL, last_batch_padded = True  -> last batch = ``[7]``,\n    next iteration will return ``[1, 2]``\n\n    last_batch_policy = LastBatchPolicy.PARTIAL, last_batch_padded = False -> last batch = ``[7]``,\n    next iteration will return ``[2, 3]``\n\n    last_batch_policy = LastBatchPolicy.FILL, last_batch_padded = True   -> last batch = ``[7, 7]``,\n    next iteration will return ``[1, 2]``\n\n    last_batch_policy = LastBatchPolicy.FILL, last_batch_padded = False  -> last batch = ``[7, 1]``,\n    next iteration will return ``[2, 3]``\n\n    last_batch_policy = LastBatchPolicy.DROP, last_batch_padded = True   -> last batch = ``[5, 6]``,\n    next iteration will return ``[1, 2]``\n\n    last_batch_policy = LastBatchPolicy.DROP, last_batch_padded = False  -> last batch = ``[5, 6]``,\n    next iteration will return ``[2, 3]``\n    \"\"\"\n\n    def __init__(self,\n                 pipelines,\n                 output_map,\n                 size=-1,\n                 reader_name=None,\n                 auto_reset=False,\n                 fill_last_batch=None,\n                 dynamic_shape=False,\n                 last_batch_padded=False,\n                 last_batch_policy=LastBatchPolicy.FILL,\n                 prepare_first_batch=True):\n\n        # check the assert first as _DaliBaseIterator would run the prefetch\n        assert len(set(output_map)) == len(output_map), \"output_map names should be distinct\"\n        self._output_categories = set(output_map)\n        self.output_map = output_map\n\n        _DaliBaseIterator.__init__(self,\n                                   pipelines,\n                                   size,\n                                   reader_name,\n                                   auto_reset,\n                                   fill_last_batch,\n                                   last_batch_padded,\n                                   last_batch_policy,\n                                   prepare_first_batch=prepare_first_batch)\n\n        self._first_batch = None\n        if self._prepare_first_batch:\n            try:\n                self._first_batch = DALIGenericIterator.__next__(self)\n                # call to `next` sets _ever_consumed to True but if we are just calling it from\n                # here we should set if to False again\n                self._ever_consumed = False\n            except StopIteration:\n                assert False, \"It seems that there is no data in the pipeline. This may happen \" \\\n                       \"if `last_batch_policy` is set to PARTIAL and the requested batch size is \" \\\n                       \"greater than the shard size.\"\n\n    def __next__(self):\n        self._ever_consumed = True\n        if self._first_batch is not None:\n            batch = self._first_batch\n            self._first_batch = None\n            return batch\n\n        # Gather outputs\n        outputs = self._get_outputs()\n\n        data_batches = [None for i in range(self._num_gpus)]\n        for i in range(self._num_gpus):\n            dev_id = self._pipes[i].device_id\n            # initialize dict for all output categories\n            category_outputs = dict()\n            # segregate outputs into categories\n            for j, out in enumerate(outputs[i]):\n                category_outputs[self.output_map[j]] = out\n\n            # Change DALI TensorLists into Tensors\n            category_tensors = dict()\n            category_shapes = dict()\n            for category, out in category_outputs.items():\n                category_tensors[category] = out.as_tensor()\n                category_shapes[category] = category_tensors[category].shape()\n\n            category_torch_type = dict()\n            category_device = dict()\n            torch_gpu_device = None\n            torch_cpu_device = torch.device('cpu')\n            # check category and device\n            for category in self._output_categories:\n                category_torch_type[category] = to_torch_type[category_tensors[category].dtype]\n                if type(category_tensors[category]) is TensorGPU:\n                    if not torch_gpu_device:\n                        torch_gpu_device = torch.device('cuda', dev_id)\n                    category_device[category] = torch_gpu_device\n                else:\n                    category_device[category] = torch_cpu_device\n\n            pyt_tensors = dict()\n            for category in self._output_categories:\n                pyt_tensors[category] = torch.empty(category_shapes[category],\n                                                    dtype=category_torch_type[category],\n                                                    device=category_device[category])\n\n            data_batches[i] = pyt_tensors\n\n            # Copy data from DALI Tensors to torch tensors\n            for category, tensor in category_tensors.items():\n                if isinstance(tensor, (TensorGPU, TensorListGPU)):\n                    # Using same cuda_stream used by torch.zeros to set the memory\n                    stream = torch.cuda.current_stream(device=pyt_tensors[category].device)\n                    feed_ndarray(tensor, pyt_tensors[category], cuda_stream=stream)\n                else:\n                    feed_ndarray(tensor, pyt_tensors[category])\n\n        self._schedule_runs()\n\n        self._advance_and_check_drop_last()\n\n        if self._reader_name:\n            if_drop, left = self._remove_padded()\n            if np.any(if_drop):\n                output = []\n                for batch, to_copy in zip(data_batches, left):\n                    batch = batch.copy()\n                    for category in self._output_categories:\n                        batch[category] = batch[category][0:to_copy]\n                    output.append(batch)\n                return output\n\n        else:\n            if self._last_batch_policy == LastBatchPolicy.PARTIAL and (\n                                          self._counter > self._size) and self._size > 0:\n                # First calculate how much data is required to return exactly self._size entries.\n                diff = self._num_gpus * self.batch_size - (self._counter - self._size)\n                # Figure out how many GPUs to grab from.\n                numGPUs_tograb = int(np.ceil(diff / self.batch_size))\n                # Figure out how many results to grab from the last GPU\n                # (as a fractional GPU batch may be required to bring us\n                # right up to self._size).\n                mod_diff = diff % self.batch_size\n                data_fromlastGPU = mod_diff if mod_diff else self.batch_size\n\n                # Grab the relevant data.\n                # 1) Grab everything from the relevant GPUs.\n                # 2) Grab the right data from the last GPU.\n                # 3) Append data together correctly and return.\n                output = data_batches[0:numGPUs_tograb]\n                output[-1] = output[-1].copy()\n                for category in self._output_categories:\n                    output[-1][category] = output[-1][category][0:data_fromlastGPU]\n                return output\n\n        return data_batches\n\n\nclass DALIClassificationIterator(DALIGenericIterator):\n    \"\"\"\n    DALI iterator for classification tasks for PyTorch. It returns 2 outputs\n    (data and label) in the form of PyTorch's Tensor.\n\n    Calling\n\n    .. code-block:: python\n\n       DALIClassificationIterator(pipelines, reader_name)\n\n    is equivalent to calling\n\n    .. code-block:: python\n\n       DALIGenericIterator(pipelines, [\"data\", \"label\"], reader_name)\n\n    Parameters\n    ----------\n    pipelines : list of nvidia.dali.Pipeline\n                List of pipelines to use\n    size : int, default = -1\n                Number of samples in the shard for the wrapped pipeline (if there is more than\n                one it is a sum)\n                Providing -1 means that the iterator will work until StopIteration is raised\n                from the inside of iter_setup(). The options `last_batch_policy` and\n                `last_batch_padded` don't work in such case. It works with only one pipeline inside\n                the iterator.\n                Mutually exclusive with `reader_name` argument\n    reader_name : str, default = None\n                Name of the reader which will be queried to the shard size, number of shards and\n                all other properties necessary to count properly the number of relevant and padded\n                samples that iterator needs to deal with. It automatically sets `last_batch_policy`\n                to PARTIAL when the FILL is used, and `last_batch_padded` accordingly to match\n                the reader's configuration\n    auto_reset : string or bool, optional, default = False\n                Whether the iterator resets itself for the next epoch or it requires reset() to be\n                called explicitly.\n\n                It can be one of the following values:\n\n                * ``\"no\"``, ``False`` or ``None`` - at the end of epoch StopIteration is raised\n                  and reset() needs to be called\n                * ``\"yes\"`` or ``\"True\"``- at the end of epoch StopIteration is raised but reset()\n                  is called internally automatically\n\n    dynamic_shape : any, optional,\n                Parameter used only for backward compatibility.\n    fill_last_batch : bool, optional, default = None\n                **Deprecated** Please use ``last_batch_policy`` instead\n\n                Whether to fill the last batch with data up to 'self.batch_size'.\n                The iterator would return the first integer multiple\n                of self._num_gpus * self.batch_size entries which exceeds 'size'.\n                Setting this flag to False will cause the iterator to return\n                exactly 'size' entries.\n    last_batch_policy: optional, default = LastBatchPolicy.FILL\n                What to do with the last batch when there are not enough samples in the epoch\n                to fully fill it. See :meth:`nvidia.dali.plugin.base_iterator.LastBatchPolicy`\n    last_batch_padded : bool, optional, default = False\n                Whether the last batch provided by DALI is padded with the last sample\n                or it just wraps up. In the conjunction with ``last_batch_policy`` it tells\n                if the iterator returning last batch with data only partially filled with\n                data from the current epoch is dropping padding samples or samples from\n                the next epoch. If set to ``False`` next\n                epoch will end sooner as data from it was consumed but dropped. If set to\n                True next epoch would be the same length as the first one. For this to happen,\n                the option `pad_last_batch` in the reader needs to be set to True as well.\n                It is overwritten when `reader_name` argument is provided\n    prepare_first_batch : bool, optional, default = True\n                Whether DALI should buffer the first batch right after the creation of the iterator,\n                so one batch is already prepared when the iterator is prompted for the data\n\n    Example\n    -------\n    With the data set ``[1,2,3,4,5,6,7]`` and the batch size 2:\n\n    last_batch_policy = LastBatchPolicy.PARTIAL, last_batch_padded = True  -> last batch = ``[7]``,\n    next iteration will return ``[1, 2]``\n\n    last_batch_policy = LastBatchPolicy.PARTIAL, last_batch_padded = False -> last batch = ``[7]``,\n    next iteration will return ``[2, 3]``\n\n    last_batch_policy = LastBatchPolicy.FILL, last_batch_padded = True   -> last batch = ``[7, 7]``,\n    next iteration will return ``[1, 2]``\n\n    last_batch_policy = LastBatchPolicy.FILL, last_batch_padded = False  -> last batch = ``[7, 1]``,\n    next iteration will return ``[2, 3]``\n\n    last_batch_policy = LastBatchPolicy.DROP, last_batch_padded = True   -> last batch = ``[5, 6]``,\n    next iteration will return ``[1, 2]``\n\n    last_batch_policy = LastBatchPolicy.DROP, last_batch_padded = False  -> last batch = ``[5, 6]``,\n    next iteration will return ``[2, 3]``\n    \"\"\"\n\n    def __init__(self,\n                 pipelines,\n                 size=-1,\n                 reader_name=None,\n                 auto_reset=False,\n                 fill_last_batch=None,\n                 dynamic_shape=False,\n                 last_batch_padded=False,\n                 last_batch_policy=LastBatchPolicy.FILL,\n                 prepare_first_batch=True):\n        super(DALIClassificationIterator, self).__init__(pipelines, [\"data\", \"label\"],\n                                                         size,\n                                                         reader_name=reader_name,\n                                                         auto_reset=auto_reset,\n                                                         fill_last_batch=fill_last_batch,\n                                                         dynamic_shape=dynamic_shape,\n                                                         last_batch_padded=last_batch_padded,\n                                                         last_batch_policy=last_batch_policy,\n                                                         prepare_first_batch=prepare_first_batch)\n\n\nclass TorchPythonFunction(ops.PythonFunctionBase):\n    schema_name = \"TorchPythonFunction\"\n    ops.register_cpu_op('TorchPythonFunction')\n    ops.register_gpu_op('TorchPythonFunction')\n\n    def _torch_stream_wrapper(self, function, *ins):\n        with torch.cuda.stream(self.stream):\n            out = function(*ins)\n        self.stream.synchronize()\n        return out\n\n    def torch_wrapper(self, batch_processing, function, device, *args):\n        func = function if device == 'cpu' else \\\n               lambda *ins: self._torch_stream_wrapper(function, *ins)\n        if batch_processing:\n            return ops.PythonFunction.function_wrapper_batch(func,\n                                                             self.num_outputs,\n                                                             torch.utils.dlpack.from_dlpack,\n                                                             torch.utils.dlpack.to_dlpack,\n                                                             *args)\n        else:\n            return ops.PythonFunction.function_wrapper_per_sample(func,\n                                                                  self.num_outputs,\n                                                                  torch_dlpack.from_dlpack,\n                                                                  torch_dlpack.to_dlpack,\n                                                                  *args)\n\n    def __call__(self, *inputs, **kwargs):\n        pipeline = Pipeline.current()\n        if pipeline is None:\n            Pipeline._raise_no_current_pipeline(\"TorchPythonFunction\")\n        if self.stream is None:\n            self.stream = torch.cuda.Stream(device=pipeline.device_id)\n        return super(TorchPythonFunction, self).__call__(*inputs, **kwargs)\n\n    def __init__(self, function, num_outputs=1, device='cpu', batch_processing=False, **kwargs):\n        self.stream = None\n        super(TorchPythonFunction, self).__init__(impl_name=\"DLTensorPythonFunctionImpl\",\n                                                  function=lambda *ins:\n                                                  self.torch_wrapper(batch_processing,\n                                                                     function, device,\n                                                                     *ins),\n                                                  num_outputs=num_outputs, device=device,\n                                                  batch_processing=batch_processing, **kwargs)\n\n\nops._wrap_op(TorchPythonFunction, \"fn\", __name__)","metadata":{"_kg_hide-input":true,"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import os\nimport sys\nimport cv2\nimport glob\nimport gdcm\nimport json\nimport shutil\nimport pydicom\nimport numpy as np\nimport pandas as pd\nimport seaborn as sns\nimport matplotlib.pyplot as plt\n\nfrom tqdm.notebook import tqdm\nfrom joblib import Parallel, delayed","metadata":{"_kg_hide-input":true,"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"## Data Preparation\n\n- I use the same strategy as in https://www.kaggle.com/code/theoviel/dicom-resized-png-jpg","metadata":{}},{"cell_type":"code","source":"IMG_PATH = \"/kaggle/input/rsna-breast-cancer-detection/test_images/\"\ntest_images = glob.glob(f\"{IMG_PATH}*/*.dcm\")\n\nif DEBUG:\n    IMG_PATH = \"/kaggle/input/rsna-breast-cancer-detection/train_images/\"\n    test_images = glob.glob(f\"{IMG_PATH}*/*.dcm\")[:1000]\n    #test_images = glob.glob(f\"{IMG_PATH}10042/*.dcm\")\n    \nprint(\"Number of images :\", len(test_images))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"SAVE_FOLDER = \"/tmp/output/\"\nSIZE = 1024\n\nos.makedirs(SAVE_FOLDER, exist_ok=True)\n\nif len(test_images) > 100:\n    N_CHUNKS = 4\nelse:\n    N_CHUNKS = 1\n\nCHUNKS = [(len(test_images) / N_CHUNKS * k, len(test_images) / N_CHUNKS * (k + 1)) for k in range(N_CHUNKS)]\nCHUNKS = np.array(CHUNKS).astype(int)\n    \nJ2K_FOLDER = \"/tmp/j2k/\"","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Process jpeg compressed dicoms on GPU\n- Convert files to j2k\n- Load j2k files, resize & scale on GPU !\n- Processing is done per batch not to run out of disk space","metadata":{}},{"cell_type":"code","source":"import torch\nimport torch.nn.functional as F\nimport nvidia.dali.fn as fn\nimport nvidia.dali.types as types\nfrom nvidia.dali import pipeline_def\nfrom nvidia.dali.types import DALIDataType\nfrom pydicom.filebase import DicomBytesIO\nfrom nvidia.dali.plugin.pytorch import feed_ndarray, to_torch_type\nimport psutil\nimport gc\n\n\ndef convert_dicom_to_j2k(file, save_folder=\"\"):\n    patient = file.split('/')[-2]\n    image = file.split('/')[-1][:-4]\n    dcmfile = pydicom.dcmread(file)\n\n    if dcmfile.file_meta.TransferSyntaxUID == '1.2.840.10008.1.2.4.90':\n        with open(file, 'rb') as fp:\n            raw = DicomBytesIO(fp.read())\n            ds = pydicom.dcmread(raw)\n        offset = ds.PixelData.find(b\"\\x00\\x00\\x00\\x0C\")  #<---- the jpeg2000 header info we're looking for\n        hackedbitstream = bytearray()\n        hackedbitstream.extend(ds.PixelData[offset:])\n        with open(save_folder + f\"{patient}_{image}.jp2\", \"wb\") as binary_file:\n            binary_file.write(hackedbitstream)\n\n            \n@pipeline_def\ndef j2k_decode_pipeline(j2kfiles):\n    jpegs, _ = fn.readers.file(files=j2kfiles)\n    images = fn.experimental.decoders.image(jpegs, device='mixed', output_type=types.ANY_DATA, dtype=DALIDataType.UINT16)\n    return images","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time \ncounter = 0\nprint(f'{counter}: mem = {psutil.virtual_memory().percent}%')\n\nfor chunk in tqdm(CHUNKS):\n    os.makedirs(J2K_FOLDER, exist_ok=True)\n\n    _ = Parallel(n_jobs=2)(\n        delayed(convert_dicom_to_j2k)(img, save_folder=J2K_FOLDER)\n        for img in test_images[chunk[0]: chunk[1]]\n    )\n    counter +=1\n    print(f'{counter}: mem = {psutil.virtual_memory().percent}%')\n    \n    j2kfiles = glob.glob(J2K_FOLDER + \"*.jp2\")\n\n    if not len(j2kfiles):\n        continue\n\n    pipe = j2k_decode_pipeline(j2kfiles, batch_size=1, num_threads=2, device_id=0, debug=True)\n    pipe.build()\n    print(f'{counter}a: mem = {psutil.virtual_memory().percent}%')\n\n    for i, f in enumerate(j2kfiles):\n        patient, image = f.split('/')[-1][:-4].split('_')\n        dicom = pydicom.dcmread(IMG_PATH + f\"{patient}/{image}.dcm\")\n\n        out = pipe.run()\n\n        # Dali -> Torch\n        img = out[0][0]\n        img_torch = torch.empty(img.shape(), dtype=torch.int16, device=\"cuda\")\n        feed_ndarray(img, img_torch, cuda_stream=torch.cuda.current_stream(device=0))\n        img = img_torch.float()\n\n        # Scale, resize, invert on GPU !\n        min_, max_ = img.min(), img.max()\n        img = (img - min_) / (max_ - min_)\n\n        if SIZE:\n            img = F.interpolate(img.view(1, 1, img.size(0), img.size(1)), (SIZE, SIZE), mode=\"bilinear\")[0, 0]\n\n        if dicom.PhotometricInterpretation == \"MONOCHROME1\":\n            img = 1 - img\n\n        # Back to CPU + SAVE\n        img = (img * 255).cpu().numpy().astype(np.uint8)\n\n        cv2.imwrite(SAVE_FOLDER + f\"{patient}_{image}.png\", img)\n\n    shutil.rmtree(J2K_FOLDER)\n\n# T4 x 2 wall time = 1 min for 1000 images\n# 1 hr 11min for 54.7k test images on 1 P100 GPU","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# Process the rest on CPU","metadata":{}},{"cell_type":"code","source":"import dicomsdl\n\ndef dicomsdl_to_numpy_image(dicom, index=0):\n    info = dicom.getPixelDataInfo()\n    dtype = info['dtype']\n    if info['SamplesPerPixel'] != 1:\n        raise RuntimeError('SamplesPerPixel != 1')\n    else:\n        shape = [info['Rows'], info['Cols']]\n    outarr = np.empty(shape, dtype=dtype)\n    dicom.copyFrameData(index, outarr)\n    return outarr\n\ndef load_img_dicomsdl(f):\n    return dicomsdl_to_numpy_image(dicomsdl.open(f))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def process(f, size=512, save_folder=\"\"):\n    patient = f.split('/')[-2]\n    image = f.split('/')[-1][:-4]\n\n    dicom = pydicom.dcmread(f)\n\n    if dicom.file_meta.TransferSyntaxUID == '1.2.840.10008.1.2.4.90':  # ALREADY PROCESSED\n        return\n\n    try:\n        img = load_img_dicomsdl(f)\n    except:\n        img = dicom.pixel_array\n\n    img = (img - img.min()) / (img.max() - img.min())\n\n    if dicom.PhotometricInterpretation == \"MONOCHROME1\":\n        img = 1 - img\n\n    img = cv2.resize(img, (size, size))\n\n    cv2.imwrite(save_folder + f\"{patient}_{image}.png\", (img * 255).astype(np.uint8))","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n_ = Parallel(n_jobs=2)(\n    delayed(process)(img, size=SIZE, save_folder=SAVE_FOLDER)\n    for img in tqdm(test_images)\n)\n\n# T4 x 2 wall time = 1 min :28 for 1k images","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print(f'mem = {psutil.virtual_memory().percent}%')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"del dicom, max_, min_, to_torch_type, DALIDataType, DicomBytesIO, test_images\ndel _\ndel feed_ndarray\ndel img\ndel CHUNKS\ndel j2kfiles\ndel pipe\ndel patient, image, out\ndel img_torch\ngc.collect()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"print(f'mem = {psutil.virtual_memory().percent}%')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%whos","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"os.listdir(SAVE_FOLDER)[:10]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# check out the PNGs output\n\ndir_path = SAVE_FOLDER\n\n# Get a list of all the files in the directory\nfiles = os.listdir(dir_path)\n\n# Calculate the number of rows and columns needed to display all the images\nn_images = min(len(files), 10) #either number of files or 10, whichever is smaller\nn_rows = int(n_images / 3) + (n_images % 3 > 0)\nn_cols = min(n_images, 3)\n\n# Create a figure and a grid of subplots\nfig, axes = plt.subplots(n_rows, n_cols, figsize=(8, 8))\n\n# Flatten the array of axes to make it easier to iterate over\naxes = axes.flatten()\n\n# Loop through the list of files\nfor ax, file in zip(axes, files):\n    # Check if the file is a png file\n    if file.endswith(\".png\"):\n        # Load the image file\n        img = cv2.imread(os.path.join(dir_path, file))\n        # Display the image\n        ax.imshow(img)\n        ax.set_title(file)\n\n# Adjust the layout of the subplots\nplt.tight_layout()\n\n# Display the plot\nplt.show()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"# PNG --> ROI.png\n","metadata":{}},{"cell_type":"code","source":"%%capture \n\n# Clone yolov5 repository\n#!git clone https://github.com/ultralytics/yolov5\n# Load trained model\n#model = torch.hub.load('./yolov5', 'custom', path='/kaggle/input/rsna-breast-cancer-detection-roi-model/rsna-roi-003.pt', source='local')\nmodel = torch.hub.load('/kaggle/input/breast-cancer-roi-brest-extractor/yolov5', 'custom', path='/kaggle/input/rsna-breast-cancer-detection-roi-model/rsna-roi-003.pt', source='local')\n\ndevice = torch.device(\"cuda\")\nmodel.to(device) ## model to GPU\n\n\n#model = torch.load('/kaggle/input/breast-cancer-roi-brest-extractor', path='/kaggle/input/rsna-breast-cancer-detection-roi-model/rsna-roi-003.pt', source='local')\n\n#import torch\n#model = torch.load('/kaggle/input/rsna-breast-cancer-detection-roi-model/rsna-roi-003.pt')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"files = os.listdir(SAVE_FOLDER)\nfull_paths = []\n\n# Use a for loop to iterate over the files and get the full path\nfor file in tqdm(files):\n  full_path = os.path.join(SAVE_FOLDER, file)\n  # Append the full path to the list\n  full_paths.append(full_path)\n\n# Print the list of full paths\nprint (f'no of files: {len(full_paths)}... Examples: {full_paths[:2]}')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\n#file_list = glob.glob(os.listdir(SAVE_FOLDER)+'*.png') #list of files created above\nfile_list = full_paths\n\nos.makedirs(\"roi_cropped\", exist_ok=True) #folder to put output from this cell\nimages = []\nerror_counter = 0\n\n#for img_file in random.sample(file_list, 25):  # it is fixed to 25 random predictions - if you want to change it remeber to change plot_roi as well\n\nfor img_file in tqdm(file_list):\n    #print(img_file)\n    # Read file from file\n    frame = cv2.imread(img_file)\n    \n    # Make prediction\n    detections = model(frame)\n    \n    # Convert results to Pandas style\n    results = detections.pandas().xyxy[0].to_dict(orient=\"records\")    \n    \n    # Plot result (in 99.99% it predicts only one instance - certainly you can assure that only best prediction \n    #is used)\n    for result in results:\n        images.append(cv2.rectangle(frame, (int(result['xmin']), int(result['ymin'])), (int(result['xmax']), int(result['ymax'])), (255,0,0), 4))\n\n    ##########\n    # new... create croopped images from test folder save to new folder\n    ##########\n    if results and isinstance(results[0], dict) and len(results[0]) >= 4:\n        b_box_values = [v for i, (k, v) in enumerate(results[0].items()) if i < 4]    \n    else:\n        # handle the case where the results list is empty or the first element is not a dictionary with at least 4 items\n        error_counter +=1\n        height, width, channels = frame.shape\n        b_box_values = [0, 0, height, width]         \n    \n\n    temp_file = cv2.imread(img_file)\n    x1, y1, x2, y2 = map(int, b_box_values)\n    cropped_file = temp_file[y1:y2, x1:x2]\n    cv2.imwrite(f'/kaggle/working/roi_cropped/{os.path.basename(img_file)}', cropped_file)\n\nprint(f'Files errorered in bbox extraction: {error_counter} total files: {len(file_list)}')","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"dir_path = \"/kaggle/working/roi_cropped\"\n\n# Get a list of all the files in the directory\nfiles = os.listdir(dir_path)\n\n# Calculate the number of rows and columns needed to display all the images\nn_images = min(len(files), 10) #either number of files or 10, whichever is smaller\nn_rows = int(n_images / 3) + (n_images % 3 > 0)\nn_cols = min(n_images, 3)\n\n# Create a figure and a grid of subplots\nfig, axes = plt.subplots(n_rows, n_cols, figsize=(8, 8))\n\n# Flatten the array of axes to make it easier to iterate over\naxes = axes.flatten()\n\n# Loop through the list of files\nfor ax, file in zip(axes, files):\n    # Check if the file is a png file\n    if file.endswith(\".png\"):\n        # Load the image file\n        img = cv2.imread(os.path.join(dir_path, file))\n        # Display the image\n        ax.imshow(img)\n        ax.set_title(file)\n\n# Adjust the layout of the subplots\nplt.tight_layout()\n\n# Display the plot\nplt.show()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"from fastai import *\nfrom fastai.vision.all import *\nimport pandas as pd\n\ntrain_df = pd.read_csv('/kaggle/input/rsna-breast-cancer-detection/train.csv')\n\nif DEBUG: info_df = pd.read_csv('/kaggle/input/rsna-breast-cancer-detection/train.csv')\nelse: info_df = pd.read_csv('/kaggle/input/rsna-breast-cancer-detection/test.csv')\n\nfn2label = {\n    \"{}_{}\".format(user, fn): cancer_or_not\n    for user, fn, cancer_or_not in zip(train_df['patient_id'], train_df['image_id'].astype('str'), train_df['cancer'])}\n\ndef label_func(path):\n    return fn2label[path.stem]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"dir_path\n#list(fn2label[path.stem].items())[:10]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"path = dir_path #'/kaggle/working/roi_cropped'\ndblock = DataBlock(\n    blocks    = (ImageBlock, CategoryBlock),\n    get_items = get_image_files,\n#    get_y = label_func,\n    item_tfms=Resize(512), #the ROI images require transformation\n    splitter = RandomSplitter()\n)\ndsets = dblock.datasets(path)\ndls = dblock.dataloaders(path)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"learner_inf = load_learner('/kaggle/input/rsna-train-and-save-models/resnet18_one_pass.pkl', cpu=False)","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"os.listdir(path)[:10]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"#Predict one at at time\n\n# cropped_imgs = os.listdir(path)[:10]\n# preds = []\n# for i in cropped_imgs:\n#     imgxx = cv2.imread(os.path.join(path, i))\n#     pred = learner_inf.predict(imgxx)\n#     preds.append(pred)\n\n# preds","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"%%time\n\n# Get images\nsorted_cropped_imgs = sorted([os.path.join(path, file) for file in os.listdir(path)])\nimages = [cv2.imread(file) for file in sorted_cropped_imgs]\n\n# pass images to fast.ai learner and get predictions\ntest_dl = learner_inf.dls.test_dl(images)\npreds_batch, _ = learner_inf.get_preds(dl=test_dl)\npredsdec, _, decoded = learner_inf.get_preds(dl=test_dl, with_decoded=True)\npredsdec[:10]\n# cpu = True 13:39s for 1k\n# cpu = false, 42s !!!","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# extract the prediction, which is the 2nd value in the tensor above\nimport numpy\narray = preds_batch.numpy()\nlist_of_lists = array.tolist()\nsecond_values = [lst[1] for lst in list_of_lists]\nsecond_values[:10]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"sorted_files = sorted(os.listdir(path))\nno_extensions = [os.path.splitext(name)[0] for name in sorted_files]\nno_extensions[:10]","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"preds_df = pd.DataFrame(data = {'concat_id':no_extensions, 'cancer':second_values})\npreds_df","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"info_df['concat_id'] = info_df['patient_id'].astype(str)+'_'+info_df['image_id'].astype(str)\ninfo_df.tail()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Merge the dataframes on the 'concat_id' column\nmerged_df = pd.merge(right=info_df, left=preds_df, on='concat_id')#, suffixes=('_test', '_preds'))\n\nif DEBUG is True:\n    merged_df['prediction_id'] = merged_df['patient_id'].astype(str)+'_'+merged_df['laterality'].astype(str)\n    merged_df.rename(columns={'cancer_y': 'cancer'}, inplace=True)\n    \nmerged_df.tail()","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"# Select only the 'prediction_id' and 'cancer' columns\nif DEBUG: resulting_df = merged_df[['prediction_id', 'cancer_x']]\nelse: resulting_df = merged_df[['prediction_id', 'cancer']]\n\nresulting_df = resulting_df.groupby('prediction_id').max()\n\nresulting_df = resulting_df.sort_index()\nresulting_df","metadata":{"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"!rm -r roi_cropped\nresulting_df.to_csv('submission.csv', index=True)\nxx = pd.read_csv('/kaggle/working/submission.csv')\nxx.tail()","metadata":{"trusted":true},"execution_count":null,"outputs":[]}]}