-
Notifications
You must be signed in to change notification settings - Fork 178
Expand file tree
/
Copy pathenqueue_task.py
More file actions
59 lines (45 loc) · 1.85 KB
/
Copy pathenqueue_task.py
File metadata and controls
59 lines (45 loc) · 1.85 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
"""Worker task for adding things to a queue."""
from multiprocessing import JoinableQueue, Process
import time
from typing import List
import boto3
class EnqueueTask:
"""A Task to send a batch of records to SQS."""
def __init__(self, messages: List[str]) -> None:
"""Initialize a Task with up to 10 SQS message entries."""
self.messages = messages
def run(self, sqs_queue: boto3.resource) -> None:
"""Send messages to SQS."""
while self.messages:
response = sqs_queue.send_messages(Entries=[
{'Id': str(i), 'MessageBody': message}
for i, message in enumerate(self.messages)
])
if not response.get('Failed'):
return
# There were some failed messages, put them back and retry in a few seconds
self.messages = [
self.messages[int(failure['Id'])]
for failure in response['Failed']
]
time.sleep(2)
class Worker(Process):
"""Worker processes consumes S3 versions from the task queue and processes them."""
def __init__(self, sqs_queue_name: str, task_queue: JoinableQueue) -> None:
"""Create a new worker process.
Args:
sqs_queue_name: Name of the target SQS queue
task_queue: Thread-safe queue of EnqueueTasks to complete
"""
super().__init__()
self._task_queue = task_queue
self._queue = boto3.resource('sqs').get_queue_by_name(QueueName=sqs_queue_name)
def run(self) -> None:
"""Consume tasks from the task queue until an empty task is found."""
while True:
task = self._task_queue.get()
if task is None:
self._task_queue.task_done()
return
task.run(self._queue)
self._task_queue.task_done()