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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 20:05:32