Consuming SQS Queues with aioboto3¶
An SQS consumer is a loop of ReceiveMessage, work, and DeleteMessage, and every one of those calls is a network round trip, so the shape of the loop decides the throughput. The less obvious failure is the visibility timeout: a message whose processing outlasts it is delivered again, and in an asyncio consumer the cause can be the connection pool rather than slow work. Measured with aioboto3 15.5 against ElasticMQ, an SQS-compatible server, through a proxy adding 30 ms per request: receiving one message per call and deleting each one ran at 14 msg/s. Receiving ten per call with delete_message_batch ran at 135 msg/s from one loop and 949 msg/s from eight. On an idle queue, short polling made 277 calls in 10 s and long polling one. With 30 receive loops sharing aioboto3's default pool of 10 connections, a heartbeat that extended visibility every second was starved of connections and ran five seconds late, and 20 messages produced 41 deliveries — 21 duplicates. With a 64-connection pool and the heartbeat, 0 duplicates. On shutdown, handing seven prefetched messages back with a zero visibility timeout made them available to another consumer in 0.03 s instead of 9.99 s. This guide builds that consumer.
Prerequisites¶
- aioboto3 or aiobotocore, and an SQS queue (ElasticMQ or LocalStack locally).
- Client and pool lifecycles, from managing aioboto3 clients without leaking connections.
- The topic overview, Cloud SDKs & Object Storage.
1. Receive and delete in batches¶
Each SQS call carries up to 10 messages. A consumer that receives one and deletes one per call spends two round trips per message; batching both sides divides that by ten:
async def consume(sqs, url: str, handle) -> None:
while True:
resp = await sqs.receive_message(QueueUrl=url, MaxNumberOfMessages=10, WaitTimeSeconds=20)
messages = resp.get("Messages", [])
if not messages:
continue
done = [m for m in messages if await handle(m)]
if done:
await sqs.delete_message_batch(QueueUrl=url, Entries=[
{"Id": str(i), "ReceiptHandle": m["ReceiptHandle"]} for i, m in enumerate(done)
])
Measured with 1,000 messages and 30 ms per request: one message per receive and one delete per message, 14 msg/s; ten per receive but one delete per message, 25 msg/s, because deletes now dominated; ten per receive with batch delete, 135 msg/s — 200 calls instead of 2,000. delete_message_batch reports failures per entry in its response's Failed list rather than raising, so check it; a message whose delete failed will come back.
Verify: for each batch received, one delete call is made, and the Failed list of its response is logged.
2. Long poll, and run several receive loops¶
With WaitTimeSeconds=0, an empty queue returns immediately and the loop asks again at once. Measured on an idle queue for 10 s: 277 receive calls with short polling, one with WaitTimeSeconds=20 — and SQS bills per request. Long polling also returns as soon as a message arrives, so it costs no latency. For throughput, run several receive loops concurrently; each one is a coroutine waiting on the network most of the time:
async def run_consumer(sqs, url, handle, loops: int = 8) -> None:
async with asyncio.TaskGroup() as tg:
for _ in range(loops):
tg.create_task(consume(sqs, url, handle))
Measured: eight loops processed 1,000 messages in 1.05 s, 949 msg/s, against 135 for one. A quick check of the emulator first is worthwhile: the same test against moto's SQS ran at about 30 calls per second whether from one loop or eight, because moto handles requests one at a time — fine for testing logic, useless for measuring a consumer. The number of loops is bounded by what the handler can process: receiving faster than you can finish only grows the number of messages whose visibility clock is running.
Verify: idle consumers make about one receive call per loop every 20 s, and throughput scales with loops until the handler is the limit.
3. Size the connection pool for pollers plus everything else¶
A long poll holds a connection for up to WaitTimeSeconds. aiobotocore's default pool is 10 connections, and every receive loop, delete and visibility change shares it. With more loops than connections, deletes and heartbeats queue behind receive calls that are simply waiting for messages:
from botocore.config import Config
RECEIVE_LOOPS = 30
WORKERS = 30
sqs_config = Config(max_pool_connections=RECEIVE_LOOPS + WORKERS + 4)
async with session.client("sqs", config=sqs_config) as sqs:
...
Measured with 30 receive loops, a 2 s visibility timeout and 3 s of work per message, the heartbeat from step 4 scheduled every second: with the default pool of 10, the heartbeat's change_message_visibility call did not get a connection until 5 s after the message was received — after the message had already been delivered to another consumer — and failed with ReceiptHandleIsInvalid; 20 messages produced 41 deliveries. With a pool of 64, the same code produced 20 deliveries and finished in 4.2 s instead of 7.3 s. Nothing in the logs said "pool exhausted"; the symptom was duplicates. A separate client for receiving, with its own pool, is an alternative that keeps acknowledgements from ever waiting on pollers.
Verify: max_pool_connections is at least the number of receive loops plus the number of concurrent deletes and heartbeats.
4. Extend visibility while work is running¶
The visibility timeout must cover processing, but processing time varies. Rather than setting a timeout long enough for the worst case — which delays redelivery after a crash by the same amount — keep it moderate and extend it from a heartbeat task while the work runs:
async def keep_invisible(sqs, url: str, receipt: str, every: float, extend: int) -> None:
while True:
await asyncio.sleep(every)
await sqs.change_message_visibility(QueueUrl=url, ReceiptHandle=receipt, VisibilityTimeout=extend)
async def handle_with_heartbeat(sqs, url, message, work, timeout: int = 30) -> bool:
heartbeat = asyncio.create_task(
keep_invisible(sqs, url, message["ReceiptHandle"], every=timeout / 3, extend=timeout))
try:
await work(message)
return True
finally:
heartbeat.cancel()
Measured in the starved configuration above, a message re-delivered while its first handler was still working left that handler with a receipt handle that ElasticMQ no longer accepted: both its heartbeat and its delete failed with ReceiptHandleIsInvalid. AWS documents that SQS may instead accept a delete with an old handle without deleting the message — either way, the work is done twice. Make handlers idempotent, as in deduplicating work across replicas, and set a redrive policy so a message that fails repeatedly moves to a dead-letter queue: with maxReceiveCount 3, a poison message was received three times — ApproximateReceiveCount 1, 2, 3 — and then appeared in the dead-letter queue.
Verify: a message whose processing takes three times the visibility timeout is delivered once, and a message that always fails ends in the dead-letter queue after maxReceiveCount receives.
5. Hand back unprocessed messages on shutdown¶
A consumer that receives ten messages and is told to stop after three has seven that it holds invisible but will never process. Return them, so other consumers can take them at once instead of after the visibility timeout:
async def release(sqs, url: str, messages: list[dict]) -> None:
if messages:
await sqs.change_message_visibility_batch(QueueUrl=url, Entries=[
{"Id": str(i), "ReceiptHandle": m["ReceiptHandle"], "VisibilityTimeout": 0}
for i, m in enumerate(messages)
])
Measured with a 10 s visibility timeout: without the release, the seven unprocessed messages reached another consumer after 9.99 s; with it, after 0.03 s. In a rolling deploy, where every consumer shuts down in turn, that difference is added to the latency of every message caught mid-batch. Stop receiving first, let in-flight handlers finish within a deadline, then release what was never started — the sequence described in handling SIGTERM in asyncio services.
Verify: after a shutdown mid-batch, unprocessed messages are received by another consumer within a second.
Verification¶
An SQS consumer is correct and efficient when:
- Receives and deletes are batched, and long polling is on.
- The connection pool covers every receive loop plus deletes and heartbeats.
- Visibility is extended while work runs, handlers are idempotent, and a dead-letter queue catches poison messages.
- Unprocessed messages are released with a zero visibility timeout on shutdown.
Diagnostic Hook: track ApproximateReceiveCount on the messages you process and alert when more than a small fraction arrive with a count above 1. Redeliveries mean visibility expired before deletion — slow handlers, a starved connection pool, or a heartbeat that is not running — and in an asyncio consumer the pool is the cause most easily missed.
Pitfalls & edge cases¶
- One message per call. Measured: 14 msg/s against 135 with batching.
- More receive loops than pool connections. Measured: heartbeats ran 5 s late, 21 duplicates.
- Short polling. Measured: 277 calls in 10 s on an idle queue.
- Holding prefetched messages through shutdown. Measured: 9.99 s before others could take them.
Frequently Asked Questions¶
How do I consume SQS messages with aioboto3?
Run several coroutines that each receive up to 10 messages with WaitTimeSeconds=20, process them, and delete them with delete_message_batch. Eight such loops handled 949 msg/s against a server with 30 ms latency in testing.
Why does my async SQS consumer process messages twice?
Visibility expired before deletion. In testing, 30 long-polling loops exhausted aiobotocore's default 10-connection pool, so visibility heartbeats ran 5 s late and 20 messages were delivered 41 times. Raise max_pool_connections and extend visibility while working.
Should I use long polling with SQS?
Yes: on an idle queue, short polling made 277 calls in 10 s and long polling made one, with no added latency when messages arrive.
What should a consumer do with prefetched messages on shutdown?
Set their visibility timeout to 0; seven released messages reached another consumer in 0.03 s instead of 9.99 s.
Related¶
- Cloud SDKs & Object Storage — up to the topic overview.
- Querying DynamoDB with aioboto3 — the next AWS service most consumers write to.
- Network I/O & Protocol Handling — the section overview.