AWS SQS and SNS Testing: LocalStack, Mocking, and End-to-End Async Flows
AWS SQS and SNS are everywhere in serverless and microservice architectures. They're also surprisingly easy to test badly. This guide covers the three layers of SQS/SNS testing: unit tests with SDK mocks, integration tests with LocalStack, and end-to-end tests for complete async flows.
The Three Testing Layers
Unit tests — test your message handler logic with mocked AWS SDK calls. Fast, no infrastructure, but won't catch IAM policy issues or serialization bugs.
Integration tests — run against LocalStack (a local AWS emulator). Test real SQS/SNS behavior without AWS costs or network calls.
End-to-end tests — run against real AWS in a test account. Most realistic, but slow and requires real AWS credentials.
Start with all three if you can afford it. If you can only do two: unit + LocalStack integration.
Setting Up LocalStack with Testcontainers
LocalStack supports a huge portion of AWS services including SQS, SNS, S3, Lambda, and DynamoDB. With Testcontainers, you get a fresh LocalStack instance per test run:
@Testcontainers
class LocalStackBaseTest {
@Container
static LocalStackContainer localstack = new LocalStackContainer(
DockerImageName.parse("localstack/localstack:3.0")
)
.withServices(Service.SQS, Service.SNS);
protected SqsClient sqsClient() {
return SqsClient.builder()
.endpointOverride(localstack.getEndpointOverride(Service.SQS))
.credentialsProvider(StaticCredentialsProvider.create(
AwsBasicCredentials.create("test", "test")
))
.region(Region.US_EAST_1)
.build();
}
protected SnsClient snsClient() {
return SnsClient.builder()
.endpointOverride(localstack.getEndpointOverride(Service.SNS))
.credentialsProvider(StaticCredentialsProvider.create(
AwsBasicCredentials.create("test", "test")
))
.region(Region.US_EAST_1)
.build();
}
}SQS Integration Tests
Basic Send and Receive
class SqsIntegrationTest extends LocalStackBaseTest {
private String queueUrl;
private SqsClient sqs;
@BeforeEach
void setup() {
sqs = sqsClient();
CreateQueueResponse queue = sqs.createQueue(
CreateQueueRequest.builder().queueName("test-queue").build()
);
queueUrl = queue.queueUrl();
}
@Test
void shouldSendAndReceiveMessage() {
sqs.sendMessage(SendMessageRequest.builder()
.queueUrl(queueUrl)
.messageBody("{\"event\":\"ORDER_PLACED\",\"orderId\":\"123\"}")
.messageAttributes(Map.of(
"eventType", MessageAttributeValue.builder()
.dataType("String")
.stringValue("ORDER_PLACED")
.build()
))
.build()
);
ReceiveMessageResponse response = sqs.receiveMessage(
ReceiveMessageRequest.builder()
.queueUrl(queueUrl)
.messageAttributeNames("All")
.maxNumberOfMessages(1)
.waitTimeSeconds(5) // Long polling
.build()
);
assertThat(response.messages()).hasSize(1);
Message msg = response.messages().get(0);
assertThat(msg.body()).contains("ORDER_PLACED");
assertThat(msg.messageAttributes()).containsKey("eventType");
}
}Testing FIFO Queues
FIFO queues guarantee ordering and deduplication — critical behavior to test:
@Test
void shouldPreserveMessageOrderInFifoQueue() {
sqs.createQueue(CreateQueueRequest.builder()
.queueName("orders.fifo")
.attributes(Map.of(
QueueAttributeName.FIFO_QUEUE, "true",
QueueAttributeName.CONTENT_BASED_DEDUPLICATION, "true"
))
.build()
);
String fifoUrl = sqs.getQueueUrl(
GetQueueUrlRequest.builder().queueName("orders.fifo").build()
).queueUrl();
// Send messages with same MessageGroupId to ensure ordering
for (int i = 1; i <= 5; i++) {
sqs.sendMessage(SendMessageRequest.builder()
.queueUrl(fifoUrl)
.messageBody("message-" + i)
.messageGroupId("order-group")
.messageDeduplicationId("dedup-" + i)
.build()
);
}
List<String> received = new ArrayList<>();
for (int i = 0; i < 5; i++) {
List<Message> messages = sqs.receiveMessage(
ReceiveMessageRequest.builder()
.queueUrl(fifoUrl)
.maxNumberOfMessages(1)
.build()
).messages();
if (!messages.isEmpty()) {
received.add(messages.get(0).body());
sqs.deleteMessage(DeleteMessageRequest.builder()
.queueUrl(fifoUrl)
.receiptHandle(messages.get(0).receiptHandle())
.build()
);
}
}
assertThat(received).containsExactly(
"message-1", "message-2", "message-3", "message-4", "message-5"
);
}Testing Dead-Letter Queues
@Test
void shouldMoveToDeadLetterAfterMaxReceives() throws Exception {
// Create DLQ first
String dlqUrl = sqs.createQueue(
CreateQueueRequest.builder().queueName("orders-dlq").build()
).queueUrl();
String dlqArn = sqs.getQueueAttributes(
GetQueueAttributesRequest.builder()
.queueUrl(dlqUrl)
.attributeNames(QueueAttributeName.QUEUE_ARN)
.build()
).attributes().get(QueueAttributeName.QUEUE_ARN);
// Create main queue with redrive policy
String redrivePolicy = String.format(
"{\"maxReceiveCount\":\"2\",\"deadLetterTargetArn\":\"%s\"}", dlqArn
);
String mainUrl = sqs.createQueue(CreateQueueRequest.builder()
.queueName("orders-main")
.attributes(Map.of(
QueueAttributeName.REDRIVE_POLICY, redrivePolicy,
QueueAttributeName.VISIBILITY_TIMEOUT, "1" // 1 second for fast testing
))
.build()
).queueUrl();
// Send a message
sqs.sendMessage(SendMessageRequest.builder()
.queueUrl(mainUrl)
.messageBody("poison-pill")
.build()
);
// Receive twice without deleting (simulating failed processing)
for (int i = 0; i < 2; i++) {
List<Message> messages = sqs.receiveMessage(
ReceiveMessageRequest.builder()
.queueUrl(mainUrl)
.maxNumberOfMessages(1)
.build()
).messages();
assertThat(messages).hasSize(1);
// Don't delete — message becomes visible again after visibility timeout
Thread.sleep(1500); // Wait for visibility timeout
}
// After maxReceiveCount, message should be in DLQ
await().atMost(Duration.ofSeconds(10)).untilAsserted(() -> {
List<Message> dlqMessages = sqs.receiveMessage(
ReceiveMessageRequest.builder()
.queueUrl(dlqUrl)
.maxNumberOfMessages(1)
.build()
).messages();
assertThat(dlqMessages).hasSize(1);
assertThat(dlqMessages.get(0).body()).isEqualTo("poison-pill");
});
}SNS Integration Tests
Fan-Out to Multiple SQS Queues
SNS fan-out is one of the most common patterns — and one of the most broken in test environments:
class SnsIntegrationTest extends LocalStackBaseTest {
@Test
void shouldFanOutToAllSubscribedQueues() {
SnsClient sns = snsClient();
SqsClient sqs = sqsClient();
// Create topic
String topicArn = sns.createTopic(
CreateTopicRequest.builder().name("order-events").build()
).topicArn();
// Create two subscriber queues
String q1Url = sqs.createQueue(
CreateQueueRequest.builder().queueName("inventory-queue").build()
).queueUrl();
String q2Url = sqs.createQueue(
CreateQueueRequest.builder().queueName("notification-queue").build()
).queueUrl();
// Get queue ARNs
String q1Arn = getQueueArn(sqs, q1Url);
String q2Arn = getQueueArn(sqs, q2Url);
// Subscribe queues to topic
sns.subscribe(SubscribeRequest.builder()
.topicArn(topicArn)
.protocol("sqs")
.endpoint(q1Arn)
.attributes(Map.of("RawMessageDelivery", "true"))
.build()
);
sns.subscribe(SubscribeRequest.builder()
.topicArn(topicArn)
.protocol("sqs")
.endpoint(q2Arn)
.attributes(Map.of("RawMessageDelivery", "true"))
.build()
);
// Publish to topic
sns.publish(PublishRequest.builder()
.topicArn(topicArn)
.message("{\"event\":\"ORDER_PLACED\",\"orderId\":\"789\"}")
.build()
);
// Both queues should receive the message
assertQueueReceivesMessage(sqs, q1Url, "ORDER_PLACED");
assertQueueReceivesMessage(sqs, q2Url, "ORDER_PLACED");
}
private void assertQueueReceivesMessage(SqsClient sqs, String queueUrl, String expected) {
await().atMost(Duration.ofSeconds(5)).untilAsserted(() -> {
List<Message> messages = sqs.receiveMessage(
ReceiveMessageRequest.builder()
.queueUrl(queueUrl)
.maxNumberOfMessages(1)
.waitTimeSeconds(1)
.build()
).messages();
assertThat(messages).hasSize(1);
assertThat(messages.get(0).body()).contains(expected);
});
}
}Testing SNS Filter Policies
Filter policies let subscribers receive only relevant messages. Test them — they're easy to misconfigure:
@Test
void shouldFilterMessagesByAttribute() {
String topicArn = sns.createTopic(
CreateTopicRequest.builder().name("orders").build()
).topicArn();
String priorityQueueUrl = sqs.createQueue(
CreateQueueRequest.builder().queueName("priority-orders").build()
).queueUrl();
String standardQueueUrl = sqs.createQueue(
CreateQueueRequest.builder().queueName("standard-orders").build()
).queueUrl();
// Subscribe priority queue with filter for tier=premium
sns.subscribe(SubscribeRequest.builder()
.topicArn(topicArn)
.protocol("sqs")
.endpoint(getQueueArn(sqs, priorityQueueUrl))
.attributes(Map.of(
"FilterPolicy", "{\"tier\":[\"premium\"]}",
"RawMessageDelivery", "true"
))
.build()
);
// Subscribe standard queue for everything
sns.subscribe(SubscribeRequest.builder()
.topicArn(topicArn)
.protocol("sqs")
.endpoint(getQueueArn(sqs, standardQueueUrl))
.attributes(Map.of("RawMessageDelivery", "true"))
.build()
);
// Publish premium order
sns.publish(PublishRequest.builder()
.topicArn(topicArn)
.message("premium order message")
.messageAttributes(Map.of(
"tier", MessageAttributeValue.builder()
.dataType("String")
.stringValue("premium")
.build()
))
.build()
);
// Publish standard order
sns.publish(PublishRequest.builder()
.topicArn(topicArn)
.message("standard order message")
.messageAttributes(Map.of(
"tier", MessageAttributeValue.builder()
.dataType("String")
.stringValue("standard")
.build()
))
.build()
);
// Priority queue: only premium message
List<Message> priorityMessages = receiveAll(sqs, priorityQueueUrl);
assertThat(priorityMessages).hasSize(1);
assertThat(priorityMessages.get(0).body()).contains("premium");
// Standard queue: both messages
List<Message> standardMessages = receiveAll(sqs, standardQueueUrl);
assertThat(standardMessages).hasSize(2);
}Unit Testing with SDK Mocks
For fast unit tests of your message handling logic:
class OrderEventHandlerTest {
private SqsClient sqsClient;
private OrderEventHandler handler;
private OrderRepository repository;
@BeforeEach
void setup() {
sqsClient = mock(SqsClient.class);
repository = mock(OrderRepository.class);
handler = new OrderEventHandler(sqsClient, repository);
}
@Test
void shouldProcessOrderPlacedEvent() {
String messageBody = "{\"eventType\":\"ORDER_PLACED\",\"orderId\":\"123\",\"amount\":99.99}";
Message message = Message.builder()
.body(messageBody)
.receiptHandle("receipt-handle-123")
.build();
when(sqsClient.receiveMessage(any(ReceiveMessageRequest.class)))
.thenReturn(ReceiveMessageResponse.builder()
.messages(List.of(message))
.build());
handler.processMessages("https://sqs.us-east-1.amazonaws.com/123/orders");
verify(repository).save(argThat(order ->
order.getOrderId().equals("123") && order.getAmount() == 99.99
));
verify(sqsClient).deleteMessage(argThat(req ->
req.receiptHandle().equals("receipt-handle-123")
));
}
@Test
void shouldNotDeleteMessageOnProcessingFailure() {
String messageBody = "invalid-json";
Message message = Message.builder()
.body(messageBody)
.receiptHandle("bad-receipt")
.build();
when(sqsClient.receiveMessage(any())).thenReturn(
ReceiveMessageResponse.builder().messages(List.of(message)).build()
);
when(repository.save(any())).thenThrow(new IllegalArgumentException("bad data"));
assertThatThrownBy(() ->
handler.processMessages("https://sqs.us-east-1.amazonaws.com/123/orders")
).isInstanceOf(MessageProcessingException.class);
// Message should NOT be deleted — let it become visible again
verify(sqsClient, never()).deleteMessage(any(DeleteMessageRequest.class));
}
}Testing AWS Lambda Triggered by SQS
When Lambda is triggered by SQS, you test it differently — with the SQSEvent input object:
class OrderLambdaHandlerTest {
private OrderLambdaHandler handler;
@BeforeEach
void setup() {
handler = new OrderLambdaHandler(mock(OrderRepository.class));
}
@Test
void shouldProcessBatchOfMessages() {
SQSEvent event = new SQSEvent();
SQSEvent.SQSMessage msg1 = createMessage("order-1", "{\"orderId\":\"1\"}");
SQSEvent.SQSMessage msg2 = createMessage("order-2", "{\"orderId\":\"2\"}");
event.setRecords(List.of(msg1, msg2));
SQSBatchResponse response = handler.handleRequest(event, mock(Context.class));
// All messages processed successfully
assertThat(response.getBatchItemFailures()).isEmpty();
}
@Test
void shouldReportPartialBatchFailures() {
SQSEvent event = new SQSEvent();
event.setRecords(List.of(
createMessage("good-msg", "{\"orderId\":\"1\"}"),
createMessage("bad-msg", "not-json")
));
SQSBatchResponse response = handler.handleRequest(event, mock(Context.class));
assertThat(response.getBatchItemFailures()).hasSize(1);
assertThat(response.getBatchItemFailures().get(0).getItemIdentifier())
.isEqualTo("bad-msg");
}
private SQSEvent.SQSMessage createMessage(String messageId, String body) {
SQSEvent.SQSMessage msg = new SQSEvent.SQSMessage();
msg.setMessageId(messageId);
msg.setBody(body);
msg.setReceiptHandle("receipt-" + messageId);
return msg;
}
}Monitoring SQS/SNS in Production
Integration tests catch configuration bugs, but production issues — message age, approximate number of messages not visible, DLQ depth — need continuous monitoring. HelpMeTest monitors your downstream services that consume from SQS/SNS, flagging processing bottlenecks before they become incidents. Try it free at helpmetest.com.
Summary
- Use LocalStack via Testcontainers for SQS/SNS integration tests — no real AWS needed
- Test FIFO ordering, DLQ redrive policies, and SNS filter policies explicitly
- Unit test message handler logic with SDK mocks for fast feedback
- For Lambda-SQS, test partial batch failure (report item failures, not full batch failure)
- Always test the unhappy path: what happens when processing fails?