SQS & SNS
SQS provides durable queues; SNS provides fan-out pub/sub. Python workers use boto3 to publish events, poll queues, and build decoupled pipelines between services and Lambda functions.
Search across all documentation pages
SQS provides durable queues; SNS provides fan-out pub/sub. Python workers use boto3 to publish events, poll queues, and build decoupled pipelines between services and Lambda functions.
import boto3
sqs = boto3.client("sqs")
queue_url = sqs.get_queue_url(QueueName="orders")["QueueUrl"]
sqs.send_message(QueueUrl=queue_url, MessageBody='{"order_id": "1"}')When to reach for this:
Publish to SNS, subscribe SQS, send/receive/delete with long polling.
import json
import boto3
sns = boto3.client("sns")
sqs = boto3.client("sqs")
TOPIC_ARN = "arn:aws:sns:us-east-1:123456789012:orders-events"
QUEUE_URL = sqs.get_queue_url(QueueName="orders-worker")["QueueUrl"]
def publish_order_event(order_id: str) -> None:
sns.publish(
TopicArn=TOPIC_ARN,
Message=json.dumps({"order_id": order_id, "type": "created"}),
)
def send_to_queue(body: dict) -> None:
sqs.send_message(QueueUrl=QUEUE_URL, MessageBody=json.dumps(body))
def receive_batch(max_messages: int = 5) -> list[dict]:
resp = sqs.receive_message(
QueueUrl=QUEUE_URL,
MaxNumberOfMessages=max_messages,
WaitTimeSeconds=20,
MessageAttributeNames=["All"],
)
return resp.get("Messages", [])
def delete_message(receipt_handle: str) -> None:
sqs.delete_message(QueueUrl=QUEUE_URL, ReceiptHandle=receipt_handle)
if __name__ == "__main__":
send_to_queue({"order_id": "42"})
for msg in receive_batch():
print(msg["Body"])
delete_message(msg["ReceiptHandle"])What this demonstrates:
publish for fan-out; SQS send_message for point-to-pointWaitTimeSeconds=20) reduces empty receives and costchange_message_visibility for long jobs| Pattern | Service |
|---|---|
| Task queue | SQS + worker |
| Broadcast | SNS topic |
| SNS → SQS fan-out | Multiple queues subscribed |
# Partial batch failure reporting for Lambda SQS event source (conceptual)
def handler(event, context):
failures = []
for record in event["Records"]:
try:
process(json.loads(record["body"]))
except Exception:
failures.append({"itemIdentifier": record["messageId"]})
return {"batchItemFailures": failures}WaitTimeSeconds 10-20.| Alternative | Use When | Don't Use When |
|---|---|---|
| Kafka/MSK | High throughput, replay log | Simple AWS-native messaging |
| Celery + Redis | Existing Python task queue | Want fully managed scaling |
| EventBridge | Event routing with rules | Simple single-queue worker |
FIFO when order and exactly-once processing matter per message group; standard for highest throughput tolerant of duplicates.
Use dedupe id in DynamoDB or DB unique constraint on business key.
SNS pushes to many subscribers; SQS buffers for workers that pull at their own rate.
CloudWatch ApproximateNumberOfMessagesVisible and age of oldest message alarms.
create_queue in bootstrap scripts or IaC; apps usually only need URLs/ARNs from config.
Use for metadata (content-type, trace id) without parsing body JSON.
moto mocks send/receive; LocalStack for integration; always test visibility timeout behavior in staging.
Up to 10 with MaxNumberOfMessages - batch processing amortizes API calls.
Subscription filter policies on message JSON attributes reduce noise to each queue.
Enable SSE on queues/topics; use KMS CMK for compliance requirements.
Stack versions: This page was written for Python 3.14.0 (stable 3.14, maintenance 3.13), FastAPI 0.115+, Django 5.2, Flask 3.1, Pydantic 2, PyTorch 2.6+, pandas 2.2+, Polars 1.x, ruff 0.9+, and uv 0.6+.
Reviewed by Chris St. John·Last updated Jul 19, 2026