使用SqsAsyncClient向AWS SQS FIFO队列发送消息丢失问题排查
AWS SQS批量发送消息大量丢失的原因分析
我使用AWS SDK for Java的SqsAsyncClient向启用了内容去重和高吞吐量的FIFO队列发送10万条消息,按每10条一组批量发送,但AWS控制台仅显示约3.7-3.8万条可用消息,超6万条丢失,相关代码如下:
@Component public class QueueMessagePublisher { private static final Logger log = LoggerFactory.getLogger(QueueMessagePublisher.class); @Value("${spring.cloud.aws.sqs.region}") private String region; @Value("${spring.cloud.aws.credentials.access-key}") private String accessKey; @Value("${spring.cloud.aws.credentials.secret-key}") private String secretKey; private SqsAsyncClient sqsAsyncClient; private final ObjectMapper objectMapper; public QueueMessagePublisher(final ObjectMapper objectMapper) { this.objectMapper = objectMapper; } @PostConstruct public void init() { this.sqsAsyncClient = SqsAsyncClient.builder() .region(Region.of(region)) .credentialsProvider(StaticCredentialsProvider.create( AwsBasicCredentials.create( accessKey, secretKey))) .build(); } public void send() throws JsonProcessingException, ExecutionException, InterruptedException { List<SendMessageBatchRequestEntry> appPushes = new ArrayList<>(); for (int i=0; i<100000; i++) { appPushes.add( SendMessageBatchRequestEntry.builder() .id(UUID.randomUUID().toString()) .messageGroupId(UUID.randomUUID().toString()) .messageDeduplicationId(UUID.randomUUID().toString()) .messageBody(objectMapper.writeValueAsString(new AppPush(String.valueOf(i)))) .build() ); } final List<List<SendMessageBatchRequestEntry>> partition = Lists.partition(appPushes, 10); final String queueUrl = sqsAsyncClient.getQueueUrl(GetQueueUrlRequest .builder() .queueName("test.fifo") .build() ).get().queueUrl(); for (List<SendMessageBatchRequestEntry> messages : partition) { SendMessageBatchRequest sendMessageBatchRequest = SendMessageBatchRequest.builder() .queueUrl(queueUrl) .entries(Collections.unmodifiableCollection(messages)) .build(); sqsAsyncClient.sendMessageBatch(sendMessageBatchRequest) .thenAccept(sendMessageBatchResponse -> { log.info("Send Message Success : {} / Send Message Fail : {}", sendMessageBatchResponse.successful().size(), sendMessageBatchResponse.failed().size()); }); } } }
核心原因分析
- 异步请求未等待完成:代码中调用
sqsAsyncClient.sendMessageBatch()后仅注册了回调,但没有等待异步操作执行完成。当send()方法执行结束后,应用上下文可能终止,大量未完成的异步请求被直接取消,消息根本没发送到SQS服务器。 - 未处理批量请求失败条目:虽然回调打印了成功/失败数量,但没有对失败的消息条目做重试处理。FIFO队列的批量请求中,部分条目可能因网络波动、限流等临时错误失败,不重试就会直接丢失。
- 消息组ID滥用触发限流:每条消息都使用随机
messageGroupId,导致SQS创建10万个独立消息组。高吞吐量FIFO队列虽支持并发,但过多消息组会触发SQS内部限流,大量请求被拒绝,而代码未处理这类错误。
修复建议
- 等待所有异步请求完成:收集所有批量请求的
CompletableFuture,用CompletableFuture.allOf().join()等待全部请求执行完毕,确保所有消息都被处理。 - 实现失败条目重试:在回调中提取失败的
SendMessageBatchRequestEntry,加入重试队列并使用指数退避策略重试,避免临时错误导致消息丢失。 - 合理设置消息组ID:根据业务逻辑分组(如按用户ID、业务类型),减少消息组数量。如果不需要严格顺序,可考虑改用标准队列。
内容的提问来源于stack exchange,提问作者hun
相关产品推荐
相关产品推荐

