Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,17 +48,27 @@ import aws_sqs_batchlib
# Receive up-to 100 messages from the given queue, polling the queue for
# up-to 15 seconds to fill the batch.
res = aws_sqs_batchlib.receive_message(
QueueUrl = "https://sqs.eu-north-1.amazonaws.com/123456789012/MyQueue",
QueueUrl="https://sqs.eu-north-1.amazonaws.com/123456789012/MyQueue",
MaxNumberOfMessages=100,
WaitTimeSeconds=15,
)

# Returns messages in the same format as boto3 / botocore SQS Client
# receive_message() method.
assert res == {
'Messages': [
{'MessageId': '[.]', 'ReceiptHandle': 'AQ[.]JA==', 'MD5OfBody': '[.]', 'Body': '[.]'},
{'MessageId': '[.]', 'ReceiptHandle': 'AQ[.]wA==', 'MD5OfBody': '[.]', 'Body': '[.]'}
"Messages": [
{
"MessageId": "[.]",
"ReceiptHandle": "AQ[.]JA==",
"MD5OfBody": "[.]",
"Body": "[.]",
},
{
"MessageId": "[.]",
"ReceiptHandle": "AQ[.]wA==",
"MD5OfBody": "[.]",
"Body": "[.]",
},
# ... up-to 100 messages
]
}
Expand Down
51 changes: 26 additions & 25 deletions aws_sqs_batchlib/aws_sqs_batchlib.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,8 @@

import time
import uuid
from typing import TYPE_CHECKING, List, Optional, Sequence, Tuple, overload
from collections.abc import Sequence
from typing import TYPE_CHECKING, Optional, overload

import boto3
import boto3.session
Expand All @@ -21,27 +22,27 @@

ReceiveMessageResultTypeDef = TypedDict(
"ReceiveMessageResultTypeDef",
{"Messages": List["MessageTypeDef"]},
{"Messages": list["MessageTypeDef"]},
)

DeleteMessageBatchResultTypeDef = TypedDict(
"DeleteMessageBatchResultTypeDef",
{
"Successful": List["DeleteMessageBatchResultEntryTypeDef"],
"Failed": List["BatchResultErrorEntryTypeDef"],
"Successful": list["DeleteMessageBatchResultEntryTypeDef"],
"Failed": list["BatchResultErrorEntryTypeDef"],
},
)

SendMessageBatchResultTypeDef = TypedDict(
"SendMessageBatchResultTypeDef",
{
"Successful": List["SendMessageBatchResultEntryTypeDef"],
"Failed": List["BatchResultErrorEntryTypeDef"],
"Successful": list["SendMessageBatchResultEntryTypeDef"],
"Failed": list["BatchResultErrorEntryTypeDef"],
},
)


def create_sqs_client(session: Optional[boto3.session.Session] = None) -> "SQSClient":
def create_sqs_client(session: boto3.session.Session | None = None) -> "SQSClient":
"""Create default SQS client.

Args:
Expand All @@ -55,7 +56,7 @@ def create_sqs_client(session: Optional[boto3.session.Session] = None) -> "SQSCl

def receive_message(
sqs_client: Optional["SQSClient"] = None,
session: Optional[boto3.session.Session] = None,
session: boto3.session.Session | None = None,
**kwargs,
) -> "ReceiveMessageResultTypeDef":
"""Receive an arbitrary number of messages from an Amazon SQS queue.
Expand Down Expand Up @@ -89,7 +90,7 @@ def receive_message(
batch_size = kwargs.get("MaxNumberOfMessages", 1)
batching_window = kwargs.get("WaitTimeSeconds", 1)

batch: List["MessageTypeDef"] = []
batch: list[MessageTypeDef] = []
start = time.time()
while time.time() - start < batching_window and len(batch) < batch_size:
kwargs["WaitTimeSeconds"] = 1
Expand All @@ -103,11 +104,11 @@ def receive_message(

def delete_message_batch(
QueueUrl: str, # pylint: disable=invalid-name
Entries: List[ # pylint: disable=invalid-name
Entries: list[ # pylint: disable=invalid-name
"DeleteMessageBatchRequestEntryTypeDef"
],
sqs_client: Optional["SQSClient"] = None,
session: Optional[boto3.session.Session] = None,
session: boto3.session.Session | None = None,
) -> "DeleteMessageBatchResultTypeDef":
"""Delete an arbitrary number of messages from an Amazon SQS queue.

Expand All @@ -128,7 +129,7 @@ def delete_message_batch(
Results similar to boto3 SQS delete_message_batch() method.
"""
sqs_client = sqs_client or create_sqs_client(session)
result: "DeleteMessageBatchResultTypeDef" = {"Successful": [], "Failed": []}
result: DeleteMessageBatchResultTypeDef = {"Successful": [], "Failed": []}

while Entries:
chunk, Entries = Entries[:10], Entries[10:]
Expand All @@ -144,11 +145,11 @@ def delete_message_batch(

def send_message_batch(
QueueUrl: str, # pylint: disable=invalid-name
Entries: List[ # pylint: disable=invalid-name
Entries: list[ # pylint: disable=invalid-name
"SendMessageBatchRequestEntryTypeDef"
],
sqs_client: Optional["SQSClient"] = None,
session: Optional[boto3.session.Session] = None,
session: boto3.session.Session | None = None,
) -> "SendMessageBatchResultTypeDef":
"""Send an arbitrary number of messages to an Amazon SQS queue.

Expand All @@ -172,7 +173,7 @@ def send_message_batch(
Results similar to boto3 SQS send_message_batch() method.
"""
sqs_client = sqs_client or create_sqs_client(session)
result: "SendMessageBatchResultTypeDef" = {"Successful": [], "Failed": []}
result: SendMessageBatchResultTypeDef = {"Successful": [], "Failed": []}

while Entries:
chunk, Entries = Entries[:10], Entries[10:]
Expand All @@ -188,21 +189,21 @@ def send_message_batch(

@overload
def _divide_failures(
failed: List["BatchResultErrorEntryTypeDef"],
failed: list["BatchResultErrorEntryTypeDef"],
entries: Sequence["SendMessageBatchRequestEntryTypeDef"],
) -> Tuple[
List["BatchResultErrorEntryTypeDef"],
List["SendMessageBatchRequestEntryTypeDef"],
) -> tuple[
list["BatchResultErrorEntryTypeDef"],
list["SendMessageBatchRequestEntryTypeDef"],
]: ... # pragma: no cover


@overload
def _divide_failures(
failed: List["BatchResultErrorEntryTypeDef"],
failed: list["BatchResultErrorEntryTypeDef"],
entries: Sequence["DeleteMessageBatchRequestEntryTypeDef"],
) -> Tuple[
List["BatchResultErrorEntryTypeDef"],
List["DeleteMessageBatchRequestEntryTypeDef"],
) -> tuple[
list["BatchResultErrorEntryTypeDef"],
list["DeleteMessageBatchRequestEntryTypeDef"],
]: ... # pragma: no cover


Expand All @@ -220,8 +221,8 @@ def _divide_failures(failed, entries):
if not failed:
return [], []

not_retryable: List["BatchResultErrorEntryTypeDef"] = []
retryable: List["SendMessageBatchRequestEntryTypeDef"] = []
not_retryable: list[BatchResultErrorEntryTypeDef] = []
retryable: list[SendMessageBatchRequestEntryTypeDef] = []

entry_map = {entry["Id"]: entry for entry in entries}
for msg in failed:
Expand Down
10 changes: 6 additions & 4 deletions benchmark/end2end.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@

import aws_sqs_batchlib

logger = logging.getLogger("test")


@contextlib.contextmanager
def stopwatch():
Expand Down Expand Up @@ -62,7 +64,7 @@ def run_iteration(args, i, sqsc):
)
send_time = get_elapsed_time()
send_per_second = len(resp["Successful"]) / send_time
logging.info(
logger.info(
"[Run=%i] Sent %i messages in %03f seconds (%i / second; %i failed)",
i,
len(resp["Successful"]),
Expand All @@ -80,7 +82,7 @@ def run_iteration(args, i, sqsc):
)
receive_time = get_elapsed_time()
receive_per_second = len(resp["Messages"]) / receive_time
logging.info(
logger.info(
"[Run=%i] Received %i messages in %03f seconds (%i / second)",
i,
len(resp["Messages"]),
Expand All @@ -99,7 +101,7 @@ def run_iteration(args, i, sqsc):
)
delete_time = get_elapsed_time()
delete_per_second = len(resp["Successful"]) / delete_time
logging.info(
logger.info(
"[Run=%i] Deleted %i messages in %03f seconds (%i / second; %i failed)",
i,
len(resp["Successful"]),
Expand Down Expand Up @@ -129,7 +131,7 @@ def main():
stats["receive"].sort()
stats["send"].sort()

logging.info("Stats: %s", json.dumps(stats, indent=2))
logger.info("Stats: %s", json.dumps(stats, indent=2))


if __name__ == "__main__":
Expand Down
4 changes: 2 additions & 2 deletions tests/test_aws_sqs_batchlib.py
Original file line number Diff line number Diff line change
Expand Up @@ -365,7 +365,7 @@ def test_delete_client_retry_failures():
{"Successful": [{"Id": "1"}, {"Id": "10"}]},
]

delete_requests = [{"Id": f"{i}", "ReceiptHandle": f"{i}"} for i in range(0, 11)]
delete_requests = [{"Id": f"{i}", "ReceiptHandle": f"{i}"} for i in range(11)]

resp = aws_sqs_batchlib.delete_message_batch(
QueueUrl=sqs_queue, Entries=delete_requests, sqs_client=client_mock
Expand Down Expand Up @@ -544,7 +544,7 @@ def test_send_retry_failures():
{"Successful": [{"Id": "1"}, {"Id": "10"}]},
]

delete_requests = [{"Id": f"{i}", "MessageBody": f"{i}"} for i in range(0, 11)]
delete_requests = [{"Id": f"{i}", "MessageBody": f"{i}"} for i in range(11)]

resp = aws_sqs_batchlib.send_message_batch(
QueueUrl=sqs_queue, Entries=delete_requests, sqs_client=client_mock
Expand Down
Loading