MCPcopy Create free account
hub / github.com/airbnb/binaryalert / _enqueue

Method _enqueue

cli/manager.py:68–109  ·  view source on GitHub ↗

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]])

Source from the content-addressed store, hash-verified

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:

Callers 4

cb_copy_allMethod · 0.95
retro_fastMethod · 0.95
retro_slowMethod · 0.95
test_enqueueMethod · 0.80

Calls 2

WorkerClass · 0.90
EnqueueTaskClass · 0.90

Tested by 1

test_enqueueMethod · 0.64