Processing SQS with AWS Lambda
Lambda's SQS event source mapping turns a queue into a serverless worker pool, and this guide configures it for production as part of Managed Cloud Queues for Background Jobs in Backend Frameworks & Worker Scaling. The defaults work for a demo; production needs deliberate choices about batch size, how many concurrent invocations one queue may drive, how the visibility timeout relates to the function timeout, and where failed messages end up.
Problem Statement
An e-commerce platform sends order events to an SQS queue, and a Lambda function writes each order into a reporting database and updates search. During a flash sale, 400,000 messages arrived in ten minutes. Lambda scaled to hundreds of concurrent invocations, the reporting database ran out of connections, most invocations failed, and the retries made it worse. Some messages were processed three times; others ended up in the DLQ even though nothing was wrong with them. You want the function to drain bursts at a rate the database can take, retry only what actually failed, and never let a healthy message reach the DLQ.
Prerequisites
- An SQS standard queue (FIFO works too, with extra ordering rules noted below) and a dead-letter queue with a redrive policy.
- A Lambda function with an execution role allowing
sqs:ReceiveMessage,sqs:DeleteMessage,sqs:GetQueueAttributes, andsqs:ChangeMessageVisibilityon the queue. - Terraform, SAM, or the CLI to manage the event source mapping.
- A handler that is idempotent — every message can be delivered more than once.
Step 1 — Understand What the Mapping Does
The event source mapping is a managed poller. It long-polls the queue, assembles batches, invokes your function synchronously, and deletes the batch's messages if the invocation succeeds. If the invocation fails, it deletes nothing; the messages become visible again when their visibility timeout expires, and each receive increments their receive count toward the DLQ threshold.
poller: ReceiveMessage (up to 10 per call, several calls per batch)
-> batch of up to BatchSize messages or MaximumBatchingWindow elapsed
-> Invoke(function, batch) [synchronous]
success -> DeleteMessageBatch(all)
error -> nothing deleted; messages reappear after VisibilityTimeout
Two defaults surprise people. The poller starts with five concurrent batches and scales up quickly when the queue is deep — up to 1,250 concurrent invocations per mapping for standard queues — limited only by the function's reserved or account concurrency. And "failure" is all-or-nothing for the batch unless you opt into partial batch responses (Step 5).
Step 2 — Set the Timeouts So They Nest Correctly
The queue's visibility timeout must be longer than the function's timeout, or messages from a still-running invocation become visible and are delivered again while the first invocation is working on them. AWS's guidance is at least six times the function timeout, which leaves room for the poller's own retries on throttling.
resource "aws_lambda_function" "orders" {
function_name = "orders-projector"
runtime = "python3.12"
handler = "handler.main"
timeout = 30 # seconds; p99 batch processing time + margin
memory_size = 512
reserved_concurrent_executions = 60 # hard ceiling for this function
# ...
}
resource "aws_sqs_queue" "orders" {
name = "orders"
visibility_timeout_seconds = 180 # >= 6 x function timeout
message_retention_seconds = 345600 # 4 days
redrive_policy = jsonencode({
deadLetterTargetArn = aws_sqs_queue.orders_dlq.arn
maxReceiveCount = 5 # see Step 6 for why not 1 or 2
})
}
The function timeout itself should cover processing a full batch, not one message. A function that handles one message in 200 ms needs about 2 seconds for a batch of 10 plus cold-start headroom; with a batch of 100 it needs 20 seconds or more.
Step 3 — Cap Concurrency at the Mapping, Not Just the Function
The flash-sale outage came from unbounded concurrency: the database allows about 200 connections, and Lambda happily ran hundreds of invocations, each opening one. The event source mapping's maximum_concurrency caps how many concurrent invocations this queue can drive, without the throttling side effects of reserved concurrency.
resource "aws_lambda_event_source_mapping" "orders" {
event_source_arn = aws_sqs_queue.orders.arn
function_name = aws_lambda_function.orders.arn
batch_size = 50 # messages per invocation
maximum_batching_window_in_seconds = 2 # wait up to 2s to fill a batch
function_response_types = ["ReportBatchItemFailures"] # Step 5
scaling_config {
maximum_concurrency = 40 # <= DB connections / conns per invocation
}
}
Why not rely on reserved_concurrent_executions alone? When a function hits its reserved concurrency, the poller's invocations are throttled, and throttled batches go back to the queue — each throttle counts as a receive and pushes healthy messages toward the DLQ. maximum_concurrency stops the poller from requesting more invocations in the first place, so no receive is wasted. Keep reserved concurrency slightly above the mapping's maximum as a backstop. The same principle — cap consumers at the downstream's capacity — is discussed in backpressure strategies for fast producers.
Step 4 — Write a Handler That Reuses Connections
Initialise clients and connections outside the handler so warm invocations reuse them; open one database connection per execution environment, not per message.
# handler.py
import json, os
import psycopg
# Created once per execution environment, reused across warm invocations
conn = psycopg.connect(os.environ["DATABASE_URL"], autocommit=False)
def process(order: dict) -> None:
with conn.transaction():
conn.execute(
"""INSERT INTO order_facts (order_id, status, total_cents, updated_at)
VALUES (%s, %s, %s, %s)
ON CONFLICT (order_id) DO UPDATE
SET status = EXCLUDED.status, total_cents = EXCLUDED.total_cents,
updated_at = EXCLUDED.updated_at
WHERE order_facts.updated_at < EXCLUDED.updated_at""", # idempotent + ordered
(order["id"], order["status"], order["total_cents"], order["updated_at"]))
def main(event, context):
failures = []
for record in event["Records"]:
try:
process(json.loads(record["body"]))
except Exception:
failures.append({"itemIdentifier": record["messageId"]})
return {"batchItemFailures": failures}
The conditional upsert makes duplicates and out-of-order redeliveries harmless — the approach from handling out-of-order events with sequence numbers. With the concurrency cap at 40 and one connection per environment, the function never opens more than about 40 connections. For higher fan-in, put RDS Proxy in front of the database.
Step 5 — Report Partial Batch Failures
With ReportBatchItemFailures enabled, the handler returns the ids of the messages that failed; the mapping deletes the rest. Without it, one bad message in a batch of 50 makes all 50 reappear and count a receive — the reason healthy messages were reaching the DLQ in the problem statement.
# Returning this deletes 48 messages and returns only 2 to the queue
{"batchItemFailures": [{"itemIdentifier": "059f36b4-..."}, {"itemIdentifier": "a1c9..."}]}
If the handler raises an exception instead of returning, the whole batch fails regardless of the setting. Catch per-record errors inside the loop, and let only truly global failures (the database is unreachable) raise. The edge cases — FIFO queues, empty lists, malformed responses — are covered in handling partial batch failures in Lambda with SQS.
Step 6 — Choose maxReceiveCount with Throttles in Mind
Every time a message is received — including when its batch was throttled or the whole invocation timed out — its receive count goes up. A maxReceiveCount of 1 or 2 will dead-letter healthy messages during any burst. Five is a sensible minimum for Lambda consumers; raise it if you see messages in the DLQ that process fine on replay.
# Messages in the DLQ that succeed when redriven were victims of throttling, not poison
aws sqs start-message-move-task --source-arn "$DLQ_ARN" --max-number-of-messages-per-second 20
Pair the DLQ with an alarm on its depth, as described in alerting on dead-letter queue growth.
Verification
Replay a burst in staging and check the three failure signals from the problem statement:
# Send 50k synthetic messages quickly
python load.py --queue orders-staging --count 50000 --rate 2000
# Concurrency never exceeds the cap
aws cloudwatch get-metric-statistics --namespace AWS/Lambda --metric-name ConcurrentExecutions \
--dimensions Name=FunctionName,Value=orders-projector --statistics Maximum --period 60 \
--start-time "$(date -u -d '-30 min' +%FT%TZ)" --end-time "$(date -u +%FT%TZ)"
# No throttles, and the DLQ stays empty
aws cloudwatch get-metric-statistics --namespace AWS/Lambda --metric-name Throttles ...
aws sqs get-queue-attributes --queue-url "$DLQ_URL" --attribute-names ApproximateNumberOfMessages
Expected: ConcurrentExecutions flat at or below 40, zero throttles, zero DLQ messages, and ApproximateAgeOfOldestMessage rising during the burst then falling back to near zero.
Gotchas & Edge Cases
FIFO queues scale differently. With FIFO, the mapping processes message groups in order and scales with the number of active groups; a single hot group runs one batch at a time. Partial batch failure responses must include the failed message and every later message in the same group.
Batching window adds latency. maximum_batching_window_in_seconds trades latency for fewer, fuller invocations. For latency-sensitive queues, keep it at 0–1 seconds.
Payload size limits. A Lambda invocation payload is limited to 6 MB; a batch of large messages can exceed it, and the mapping reduces the batch automatically. Keep messages small.
Cold starts in bursts. Scaling from 5 to 40 concurrent invocations creates new execution environments, each paying initialisation time and opening a database connection. Provisioned concurrency smooths this for latency-critical queues.
FAQ
Should I use Lambda or containers to consume SQS? Lambda for short, spiky, independent jobs where zero idle cost matters. Containers for long jobs (over a few minutes), steady high throughput where per-invocation pricing adds up, or jobs needing large memory, GPUs, or persistent connections.
Why do I see duplicate processing even without errors? A visibility timeout shorter than processing time, or SQS standard's at-least-once delivery. Check that the visibility timeout is at least six times the function timeout, and make the handler idempotent.
Can one queue trigger several functions? Several mappings on one queue compete for messages — each message goes to one of them. For fan-out, publish to SNS or EventBridge with one queue per consumer.
Related
- Managed Cloud Queues for Background Jobs — delivery models and provider trade-offs.
- Handling Partial Batch Failures in Lambda with SQS — the response contract in detail.
- Configuring an SQS Redrive Policy — DLQ setup and replay.
- Debugging Duplicate Deliveries in SQS — tracing the causes of repeats.