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, and sqs:ChangeMessageVisibility on 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).

A managed poller between queue and function The event source mapping receives messages from the SQS queue and assembles a batch. It invokes the Lambda function synchronously with the batch. If the function succeeds, the mapping deletes all messages in the batch. If it fails, nothing is deleted and the messages become visible again after the visibility timeout, counting toward the DLQ threshold. Poll, batch, invoke, delete SQS queue orders event source mapping BatchSize, window, max conc. Lambda function handler(batch) success: DeleteMessageBatch error: nothing deleted, batch reappears

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
  })
}
Timeouts must nest The p99 time to process one batch is about 12 seconds. The function timeout of 30 seconds contains it with margin. The queue visibility timeout of 180 seconds contains the function timeout six times over, so a message is never redelivered while an invocation holding it is still running, even after poller retries. batch p99 < function timeout < visibility timeout visibility timeout: 180 s function timeout 30 s batch 12 s room for poller retries and throttled re-invocations If the inner box ever outgrows the outer one, the same batch runs twice at once.

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.

Cap the poller at the database's capacity Without a cap, a deep queue drives hundreds of concurrent invocations, each opening a database connection, exceeding the database's 200-connection limit and causing failures and throttle-driven redeliveries. With maximum concurrency set to 40 on the event source mapping, at most 40 invocations run, the database stays within its limit, and the backlog drains steadily. Concurrent invocations during a 400k-message burst uncapped ~600 invocations: DB connections exhausted max_concurrency 40 40 invocations, backlog drains in ~25 min DB limit: 200 connections A slower drain that succeeds beats a fast one that fails and retries.

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