A Task to send a batch of records to SQS.
| 7 | |
| 8 | |
| 9 | class EnqueueTask: |
| 10 | """A Task to send a batch of records to SQS.""" |
| 11 | |
| 12 | def __init__(self, messages: List[str]) -> None: |
| 13 | """Initialize a Task with up to 10 SQS message entries.""" |
| 14 | self.messages = messages |
| 15 | |
| 16 | def run(self, sqs_queue: boto3.resource) -> None: |
| 17 | """Send messages to SQS.""" |
| 18 | while self.messages: |
| 19 | response = sqs_queue.send_messages(Entries=[ |
| 20 | {'Id': str(i), 'MessageBody': message} |
| 21 | for i, message in enumerate(self.messages) |
| 22 | ]) |
| 23 | |
| 24 | if not response.get('Failed'): |
| 25 | return |
| 26 | |
| 27 | # There were some failed messages, put them back and retry in a few seconds |
| 28 | self.messages = [ |
| 29 | self.messages[int(failure['Id'])] |
| 30 | for failure in response['Failed'] |
| 31 | ] |
| 32 | time.sleep(2) |
| 33 | |
| 34 | |
| 35 | class Worker(Process): |