# Program to demonstrate using thread-pools to improve # throughput on blocking I/O calls. # # Sample output (4 workers, 50 messages): # 2021-06-17 15:42:37,673 sending msgs: starting # 2021-06-17 15:42:37,675 sending msg: 0 # 2021-06-17 15:42:37,675 call to send msg 1 / 50 # 2021-06-17 15:42:37,675 sending msg: 1 # 2021-06-17 15:42:37,675 call to send msg 2 / 50 # 2021-06-17 15:42:37,675 sending msg: 2 # 2021-06-17 15:42:37,675 call to send msg 3 / 50 # 2021-06-17 15:42:37,676 sending msg: 3 # 2021-06-17 15:42:37,676 call to send msg 4 / 50 # 2021-06-17 15:42:37,676 call to send msg 5 / 50 # 2021-06-17 15:42:37,676 call to send msg 6 / 50 # 2021-06-17 15:42:37,676 call to send msg 7 / 50 # 2021-06-17 15:42:37,676 call to send msg 8 / 50 # 2021-06-17 15:42:37,676 call to send msg 9 / 50 # 2021-06-17 15:42:37,676 call to send msg 10 / 50 # 2021-06-17 15:42:37,676 call to send msg 11 / 50 # 2021-06-17 15:42:37,676 call to send msg 12 / 50 # 2021-06-17 15:42:37,676 call to send msg 13 / 50 # 2021-06-17 15:42:37,676 call to send msg 14 / 50 # 2021-06-17 15:42:37,676 call to send msg 15 / 50 # 2021-06-17 15:42:37,677 call to send msg 16 / 50 # 2021-06-17 15:42:37,677 call to send msg 17 / 50 # 2021-06-17 15:42:37,677 call to send msg 18 / 50 # 2021-06-17 15:42:37,677 call to send msg 19 / 50 # 2021-06-17 15:42:37,677 call to send msg 20 / 50 # 2021-06-17 15:42:37,677 call to send msg 21 / 50 # 2021-06-17 15:42:37,677 call to send msg 22 / 50 # 2021-06-17 15:42:37,677 call to send msg 23 / 50 # 2021-06-17 15:42:37,677 call to send msg 24 / 50 # 2021-06-17 15:42:37,677 call to send msg 25 / 50 # 2021-06-17 15:42:37,677 call to send msg 26 / 50 # 2021-06-17 15:42:37,677 call to send msg 27 / 50 # 2021-06-17 15:42:37,677 call to send msg 28 / 50 # 2021-06-17 15:42:37,677 call to send msg 29 / 50 # 2021-06-17 15:42:37,677 call to send msg 30 / 50 # 2021-06-17 15:42:37,677 call to send msg 31 / 50 # 2021-06-17 15:42:37,677 call to send msg 32 / 50 # 2021-06-17 15:42:37,677 call to send msg 33 / 50 # 2021-06-17 15:42:37,677 call to send msg 34 / 50 # 2021-06-17 15:42:37,677 call to send msg 35 / 50 # 2021-06-17 15:42:37,677 call to send msg 36 / 50 # 2021-06-17 15:42:37,677 call to send msg 37 / 50 # 2021-06-17 15:42:37,677 call to send msg 38 / 50 # 2021-06-17 15:42:37,678 call to send msg 39 / 50 # 2021-06-17 15:42:37,678 call to send msg 40 / 50 # 2021-06-17 15:42:37,678 call to send msg 41 / 50 # 2021-06-17 15:42:37,678 call to send msg 42 / 50 # 2021-06-17 15:42:37,678 call to send msg 43 / 50 # 2021-06-17 15:42:37,678 call to send msg 44 / 50 # 2021-06-17 15:42:37,678 call to send msg 45 / 50 # 2021-06-17 15:42:37,678 call to send msg 46 / 50 # 2021-06-17 15:42:37,678 call to send msg 47 / 50 # 2021-06-17 15:42:37,678 call to send msg 48 / 50 # 2021-06-17 15:42:37,678 call to send msg 49 / 50 # 2021-06-17 15:42:37,678 call to send msg 50 / 50 # 2021-06-17 15:42:38,677 sending msg: 4 # 2021-06-17 15:42:38,677 sending msg: 5 # 2021-06-17 15:42:38,678 sending msg: 6 # 2021-06-17 15:42:38,678 sending msg failed for msg 3: send failed on msg: 3 # 2021-06-17 15:42:38,678 sending msg: 7 # 2021-06-17 15:42:39,681 sending msg: 8 # 2021-06-17 15:42:39,681 sending msg: 9 # 2021-06-17 15:42:39,681 sending msg: 10 # 2021-06-17 15:42:39,681 sending msg: 11 # 2021-06-17 15:42:39,681 sending msg failed for msg 7: send failed on msg: 7 # 2021-06-17 15:42:40,686 sending msg: 12 # 2021-06-17 15:42:40,686 sending msg: 13 # 2021-06-17 15:42:40,687 sending msg: 14 # 2021-06-17 15:42:40,687 sending msg: 15 # 2021-06-17 15:42:41,689 sending msg: 16 # 2021-06-17 15:42:41,689 sending msg: 17 # 2021-06-17 15:42:41,690 sending msg: 18 # 2021-06-17 15:42:41,690 sending msg: 19 # 2021-06-17 15:42:42,692 sending msg: 20 # 2021-06-17 15:42:42,692 sending msg: 21 # 2021-06-17 15:42:42,692 sending msg: 22 # 2021-06-17 15:42:42,692 sending msg: 23 # 2021-06-17 15:42:43,694 sending msg: 24 # 2021-06-17 15:42:43,694 sending msg: 25 # 2021-06-17 15:42:43,694 sending msg: 26 # 2021-06-17 15:42:43,695 sending msg: 27 # 2021-06-17 15:42:44,695 sending msg: 28 # 2021-06-17 15:42:44,695 sending msg: 29 # 2021-06-17 15:42:44,700 sending msg: 30 # 2021-06-17 15:42:44,700 sending msg: 31 # 2021-06-17 15:42:45,699 sending msg: 32 # 2021-06-17 15:42:45,700 sending msg: 33 # 2021-06-17 15:42:45,700 sending msg: 34 # 2021-06-17 15:42:45,705 sending msg: 35 # 2021-06-17 15:42:46,705 sending msg: 36 # 2021-06-17 15:42:46,705 sending msg: 37 # 2021-06-17 15:42:46,705 sending msg: 38 # 2021-06-17 15:42:46,706 sending msg: 39 # 2021-06-17 15:42:47,709 sending msg: 40 # 2021-06-17 15:42:47,709 sending msg: 41 # 2021-06-17 15:42:47,710 sending msg: 42 # 2021-06-17 15:42:47,710 sending msg: 43 # 2021-06-17 15:42:48,712 sending msg: 44 # 2021-06-17 15:42:48,712 sending msg: 45 # 2021-06-17 15:42:48,713 sending msg: 46 # 2021-06-17 15:42:48,713 sending msg: 47 # 2021-06-17 15:42:49,715 sending msg: 48 # 2021-06-17 15:42:49,715 sending msg: 49 # 2021-06-17 15:42:50,719 sending msgs: finished import logging import time import concurrent.futures import random num_msgs = 50 num_workers = 4 failure_pct = 10 def send_message(msg): logging.info('sending msg: %s', msg) # this simulates slowness time.sleep(1) if random.randint(0, 99) % failure_pct == msg: # this simulates failure raise IOError(f'send failed on msg: {msg}') def send_messages_quickly_method2(msgs): logging.info('sending msgs: starting') with concurrent.futures.ThreadPoolExecutor(max_workers=num_workers) as pool: futures = {} # keyed by future obj, value can be something we need for reference later # Submit our tasks to the worker pool for i, msg in enumerate(msgs): future = pool.submit(send_message, msg) logging.info('call to send msg %d / %d', i + 1, len(msgs)) # we keep a reference to our message for use later futures[future] = msg # Wait for workers to finish the tasks for future in concurrent.futures.as_completed(futures): msg = None try: msg = futures[future] _ = future.result() # raises an exception if post failed except IOError as e: # match this with whatever exception your I/O API raises logging.warning('sending msg failed for msg %s: %s', msg, e) logging.info('sending msgs: finished') if __name__ == '__main__': logging.basicConfig(format='%(asctime)s %(message)s', level=logging.INFO) random.seed(42) msgs = list(range(num_msgs)) send_messages_quickly_method2(msgs)