You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot中SQS消费者无故停止消费问题排查求助

Spring Boot SQS Listener Stops Consuming Messages Without Errors (ApproximateNumberOfMessagesVisible Rising)

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 declares throws IOException, but if indexerService.indexSet(setName) throws an unchecked exception (like a RuntimeException subclass) 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 time
    
  • Silent connection drops in the SQS client
    If the underlying AmazonSQSAsync client 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.09 18:37:29