Spring Boot中SQS消费者无故停止消费问题排查求助
Problem Description
I'm using spring-cloud-starter-aws-messaging in a Spring Boot app to consume SQS messages via the @SqsListener annotation. Out of nowhere, the consumer stops receiving messages, the ApproximateNumberOfMessagesVisible metric keeps rising (triggering CloudWatch alarms), and there are no error logs generated before it stops.
Here's my consumer code:
@SqsListener(value = "${sqs.queue.url.indexSavedSetQueue}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS) public void listenIndexSavedSetEvent(@NonNull String message) throws IOException { log.info("Index saved set event received, message: {}", message); IndexSavedSetPayloadDto indexSavedSetPayloadDto = objectMapper .readValue(message, IndexSavedSetPayloadDto.class); String setName = indexSavedSetPayloadDto.getSetName(); indexerService.indexSet(setName); }
Possible Causes & Fixes
I've dealt with similar issues before, so here are the most likely culprits and how to fix them:
Uncaught exceptions draining the listener thread pool
Your method declaresthrows IOException, but ifindexerService.indexSet(setName)throws an unchecked exception (like aRuntimeExceptionsubclass) that you don't catch, it will kill the listener thread. Spring Cloud AWS uses a limited-size thread pool for SQS listeners—once enough threads die from uncaught exceptions, there are no threads left to process new messages.
Fix: Wrap your processing logic in a try-catch block to handle all exceptions, ensuring none escape and take down threads. Example:@SqsListener(value = "${sqs.queue.url.indexSavedSetQueue}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS) public void listenIndexSavedSetEvent(@NonNull String message) { log.info("Index saved set event received, message: {}", message); try { IndexSavedSetPayloadDto indexSavedSetPayloadDto = objectMapper .readValue(message, IndexSavedSetPayloadDto.class); String setName = indexSavedSetPayloadDto.getSetName(); indexerService.indexSet(setName); } catch (Exception e) { log.error("Failed to process index saved set event", e); // Optional: Throw a wrapped exception to trigger retry or DLQ handling throw new MessageConversionException("Failed to process message", e); } }Insufficient default thread pool configuration
The default thread pool size for Spring Cloud AWS SQS listeners is small. If your message volume spikes or each message takes a long time to process, the pool gets fully occupied, and no new messages can be picked up.
Fix: Customize the listener container factory to adjust thread pool and message fetch settings:@Bean public SimpleMessageListenerContainerFactory simpleMessageListenerContainerFactory(AmazonSQSAsync amazonSqs) { SimpleMessageListenerContainerFactory factory = new SimpleMessageListenerContainerFactory(); factory.setAmazonSqs(amazonSqs); factory.setMaxNumberOfMessages(10); // Number of messages to fetch per poll factory.setConcurrency("5-10"); // Min/max thread pool size return factory; }Misconfigured Visibility Timeout
If your message processing time exceeds the queue's Visibility Timeout, the message will become visible again in the queue. This leads to repeated processing of the same message, which can make it look like the listener has stopped (when it's actually stuck on retries).
Fix: Ensure the queue's Visibility Timeout is longer than your average message processing time. You can also override it per listener:@SqsListener(value = "${sqs.queue.url.indexSavedSetQueue}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS, visibilityTimeout = 60) // In seconds, adjust based on your processing timeSilent connection drops in the SQS client
If the underlyingAmazonSQSAsyncclient loses connection to AWS (due to network blips or temporary service outages) and isn't configured to reconnect automatically, the listener won't resume fetching messages.
Fix: Configure the SQS client with retry and timeout settings to handle connection issues:@Bean public AmazonSQSAsync amazonSqsAsync() { ClientConfiguration clientConfig = new ClientConfiguration(); clientConfig.setMaxErrorRetry(3); // Retry failed requests clientConfig.setConnectionTimeout(5000); // 5-second connection timeout clientConfig.setSocketTimeout(5000); // 5-second socket timeout return AmazonSQSAsyncClientBuilder.standard() .withClientConfiguration(clientConfig) .withCredentials(new DefaultAWSCredentialsProviderChain()) .build(); }Missing retry or Dead-Letter Queue (DLQ) setup
Without a retry mechanism or DLQ, failed messages keep getting reprocessed, hogging thread resources and preventing new messages from being consumed.
Fix: Set up a DLQ for your SQS queue to move messages that fail processing multiple times. You can also add a retry template to handle transient failures:@Bean public RetryTemplate retryTemplate() { RetryTemplate retryTemplate = new RetryTemplate(); SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); retryPolicy.setMaxAttempts(3); // Retry up to 3 times retryTemplate.setRetryPolicy(retryPolicy); return retryTemplate; } @SqsListener(value = "${sqs.queue.url.indexSavedSetQueue}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS) public void listenIndexSavedSetEvent(@NonNull String message) { retryTemplate.execute(context -> { log.info("Index saved set event received, message: {}", message); IndexSavedSetPayloadDto indexSavedSetPayloadDto = objectMapper .readValue(message, IndexSavedSetPayloadDto.class); String setName = indexSavedSetPayloadDto.getSetName(); indexerService.indexSet(setName); return null; }); }
内容的提问来源于stack exchange,提问作者sarat chandra

