{"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":"### **This is demo of processing dicom files in parallel using Ray.** \n\nI decided to show (maybe somebody will use this idea in competition) how to process competition data in parallel.\n\nDefined tasks:\n- task1 - process dicom files in parallel\n- task2 - do something with files from task1 (*.png files) in parallel (this is sample so I decided to count files in directry byt you can do something else)\n\nFor communication purposes (between task) I used messages.\n\nRay has a lot of possibilities to process data in ML pipeline. Look here: https://docs.ray.io/en/latest/index.html#","metadata":{"_uuid":"8f2839f25d086af736a60e9eeb907d3b93b6e0e5","_cell_guid":"b1076dfc-b9ad-4769-8c92-a6c4dae69d19"}},{"cell_type":"code","source":"%%capture\n\n!pip install /kaggle/input/rsnamodules/dicomsdl-0.109.1-cp37-cp37m-manylinux_2_12_x86_64.manylinux2010_x86_64.whl ","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:01:21.28639Z","iopub.execute_input":"2022-12-09T10:01:21.287303Z","iopub.status.idle":"2022-12-09T10:01:32.510289Z","shell.execute_reply.started":"2022-12-09T10:01:21.287262Z","shell.execute_reply":"2022-12-09T10:01:32.508314Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import dicomsdl as dicoml\nimport cv2\n\nimport glob\nimport time\nimport numpy as np\nimport random\n\nfrom matplotlib import pyplot as plt","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:01:32.515184Z","iopub.execute_input":"2022-12-09T10:01:32.5159Z","iopub.status.idle":"2022-12-09T10:01:32.524163Z","shell.execute_reply.started":"2022-12-09T10:01:32.515836Z","shell.execute_reply":"2022-12-09T10:01:32.521807Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"FILE_SIZE = 512\n\n# take images patients directory\ntrain_group = glob.glob(\"/kaggle/input/rsna-breast-cancer-detection/train_images/*/\")\nprint(f'Patients for processing: {len(train_group)}')","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:01:32.526933Z","iopub.execute_input":"2022-12-09T10:01:32.527635Z","iopub.status.idle":"2022-12-09T10:01:35.784203Z","shell.execute_reply.started":"2022-12-09T10:01:32.527566Z","shell.execute_reply":"2022-12-09T10:01:35.782476Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"import ray\nray.init(log_to_driver=False)","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:01:35.787622Z","iopub.execute_input":"2022-12-09T10:01:35.788615Z","iopub.status.idle":"2022-12-09T10:01:40.645509Z","shell.execute_reply.started":"2022-12-09T10:01:35.788562Z","shell.execute_reply":"2022-12-09T10:01:40.642393Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"def process_faster(f, size=512, save_folder=None, dicom_process = True, extension=\"png\", group_id = 0):\n    \n    patient = f.split('/')[-2]\n    image_name = f.split('/')[-1][:-4]\n\n    dicom = dicoml.open(f)\n    img = dicom.pixelData()\n\n    img = (img - img.min()) / (img.max() - img.min())\n\n    if dicom.getPixelDataInfo()['PhotometricInterpretation'] == \"MONOCHROME1\":\n        img = 1 - img\n\n    image = (img * 255).astype(np.uint8)\n        \n    img = cv2.resize(image, (size, size))\n\n    file_name = f'{save_folder}' + f\"{patient}_{image_name}.{extension}\"\n    cv2.imwrite(file_name, img)\n\n# sample function - draw circle in center of the image\ndef draw_circle(img_files):\n    for f_img in img_files:\n        img = cv2.imread(f_img)\n        cv2.circle(img, (int(FILE_SIZE // 2),int(FILE_SIZE // 2)), 20, (255,0,0), 3)\n        cv2.imwrite(f_img, img)\n\n\n@ray.remote\nclass MessageActor(object):\n    def __init__(self):\n        self.messages = []\n    \n    def add_message(self, message):\n        self.messages.append(message)\n    \n    def get_and_clear_messages(self):\n        messages = self.messages\n        self.messages = []\n        return messages\n\n@ray.remote\ndef group_processor(group_id):\n    # This function process all dicom files for particular patient ID\n    files = glob.glob(f'/kaggle/input/rsna-breast-cancer-detection/train_images/{group_id}/*.dcm')\n    \n    # process all files in directory\n    for f in files:\n        res = process_faster(f, size = FILE_SIZE, save_folder = '', dicom_process = False)\n    \n    # when finished task send message to queue - \"task finished all files for patient ID were processed\" \n    message_actor.add_message.remote(f\"p-{group_id}\")\n    \n@ray.remote\ndef process_group(group_id):\n    # task X on all files from group e.g. make prediction\n    # this is sample so I decided to:\n    # a. check number of files in directory\n    # b. draw a small circle on image \n    \n    files = glob.glob(f'{group_id}_*.png')\n    count = len(files)\n    draw_circle(files)\n    # send message to queue - task finished\n    message_actor.add_message.remote(f\"d-{group_id},{count}\")\n    \nmessage_actor = MessageActor.remote()","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:01:40.650461Z","iopub.execute_input":"2022-12-09T10:01:40.651492Z","iopub.status.idle":"2022-12-09T10:01:40.684413Z","shell.execute_reply.started":"2022-12-09T10:01:40.651425Z","shell.execute_reply":"2022-12-09T10:01:40.682877Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"code","source":"LIMIT = 25 # we process only 25 patients\nall_files = 0\nstart_time = time.time()\n\n# As a sample process 20 clients\nfor g in train_group[:LIMIT]:\n    group_id = g.split('/')[-2]\n    # run it parallel\n    group_processor.remote(group_id)\n\n# Now we will listen all messages \nwhile True:\n    # waiting for task finished\n    # we have 2 task defined\n    # task1 - processing dicom files for particular patient\n    # task2 - sample function \n    \n    new_messages = ray.get(message_actor.get_and_clear_messages.remote())\n    \n    # check messages from queue\n    for message in new_messages:\n        # tokenize message\n        message_tokens = message.split('-')[-1].split(',')\n        message_type = message[0]\n        \n        # message from task1?\n        if message_type=='p':\n            print(f'Finished task1 for patient {message_tokens[0]}')\n            \n            # call second step of procesing - we know that all files for pateint X are processed \n            process_group.remote(message.split('-')[-1])\n            LIMIT -= len(new_messages)\n        \n        #message from task2?\n        elif message_type=='d':\n            all_files += int(message_tokens[-1])\n            print(f\"Finished task2 for patient {message_tokens[0]} number of files {message_tokens[-1]}\")\n    \n    # we run LIMIT task - so we check if all are processed\n    if not LIMIT:\n        break\n        \nprint(f\"End of processing {all_files} files - time: {time.time()-start_time}\")\n\nray.shutdown()","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:01:40.689861Z","iopub.execute_input":"2022-12-09T10:01:40.690604Z","iopub.status.idle":"2022-12-09T10:02:29.074115Z","shell.execute_reply.started":"2022-12-09T10:01:40.69054Z","shell.execute_reply":"2022-12-09T10:02:29.072473Z"},"trusted":true},"execution_count":null,"outputs":[]},{"cell_type":"markdown","source":"### Results:\n- Parallel processing\n- **Look on processing time!**","metadata":{}},{"cell_type":"code","source":"n = 4\nout_files = glob.glob(f'./*.png')\n\nfig, axs = plt.subplots(1, n, figsize=(25, 25))\nfor idx, im_file in enumerate(random.sample(out_files, n)):\n    im = cv2.imread(im_file)\n    axs[idx].imshow(im)\nplt.show() ","metadata":{"execution":{"iopub.status.busy":"2022-12-09T10:02:29.076572Z","iopub.execute_input":"2022-12-09T10:02:29.077232Z","iopub.status.idle":"2022-12-09T10:02:29.910398Z","shell.execute_reply.started":"2022-12-09T10:02:29.077172Z","shell.execute_reply":"2022-12-09T10:02:29.908934Z"},"trusted":true},"execution_count":null,"outputs":[]}]}