Use multiple worker processes to enqueue messages onto an SQS queue in batches. Args: queue_name: Name of the target SQS queue messages: Iterable of dictionaries, each representing a single SQS message body summary_func: Function from message to (item_cou
(
queue_name: str, messages: Iterable[Dict[str, Any]],
summary_func: Callable[[Dict[str, Any]], Tuple[int, str]])
| 66 | |
| 67 | @staticmethod |
| 68 | def _enqueue( |
| 69 | queue_name: str, messages: Iterable[Dict[str, Any]], |
| 70 | summary_func: Callable[[Dict[str, Any]], Tuple[int, str]]) -> None: |
| 71 | """Use multiple worker processes to enqueue messages onto an SQS queue in batches. |
| 72 | |
| 73 | Args: |
| 74 | queue_name: Name of the target SQS queue |
| 75 | messages: Iterable of dictionaries, each representing a single SQS message body |
| 76 | summary_func: Function from message to (item_count, summary) to show progress |
| 77 | """ |
| 78 | num_workers = multiprocessing.cpu_count() * 4 |
| 79 | tasks: JoinableQueue = JoinableQueue(num_workers * 10) # Max tasks waiting in queue |
| 80 | |
| 81 | # Create and start worker processes |
| 82 | workers = [Worker(queue_name, tasks) for _ in range(num_workers)] |
| 83 | for worker in workers: |
| 84 | worker.start() |
| 85 | |
| 86 | # Create an EnqueueTask for each batch of 10 messages (max allowed by SQS) |
| 87 | message_batch = [] |
| 88 | progress = 0 # Total number of relevant "items" processed so far |
| 89 | for message_body in messages: |
| 90 | count, summary = summary_func(message_body) |
| 91 | progress += count |
| 92 | print('\r{}: {:<90}'.format(progress, summary), end='', flush=True) |
| 93 | |
| 94 | message_batch.append(json.dumps(message_body, separators=(',', ':'))) |
| 95 | |
| 96 | if len(message_batch) == 10: |
| 97 | tasks.put(EnqueueTask(message_batch)) |
| 98 | message_batch = [] |
| 99 | |
| 100 | # Add final batch of messages |
| 101 | if message_batch: |
| 102 | tasks.put(EnqueueTask(message_batch)) |
| 103 | |
| 104 | # Add "poison pill" to mark the end of the task queue |
| 105 | for _ in range(num_workers): |
| 106 | tasks.put(None) |
| 107 | |
| 108 | tasks.join() |
| 109 | print('\nDone!') |
| 110 | |
| 111 | @staticmethod |
| 112 | def apply() -> None: |