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

为何启用并行后SQS消息消费仍过慢?求专家指导

SQS消费速度不足问题排查与优化建议

问题背景

SQS队列每15分钟接收10万条消息,基于Spring Boot构建的消费者每次拉取10条消息并行处理。当前部署10台4核8GB配置的ECS容器,但消费速度仍无法匹配消息流入速度,出现堆积。

相关代码实现

调度器代码

@Component
@Profile("!test")
public class StockRangeConsumerScheduler {

    private final StockRangeConsumer stockRangeSQSConsumer;

    public StockRangeConsumerScheduler(StockRangeConsumer stockRangeSQSConsumer) {
        this.stockRangeSQSConsumer = stockRangeSQSConsumer;
    }

    @Scheduled(fixedDelayString = "${stockrange.scheduled.delay.fixed}", initialDelayString = "${stockrange.scheduled.delay.initial}")
    public void process() throws ExecutionException, InterruptedException {
        stockRangeSQSConsumer.consume();
    }

}

消费者核心代码

@Service
@Profile({"aws", "vanguard"})
public class StockRangeConsumer {

    private static final Logger LOGGER = LoggerFactory.getLogger(StockRangeConsumer.class);

    public static final int THREADS = 4;

    private final ObjectMapper objectMapper;
    private final SqsMessageConsumer stockRangeConsumer;
    private final ForkJoinPool forkJoinPool = new ForkJoinPool(THREADS);

    private final StockRangeDataProcessor stockRangeDataProcessor;
    private final PilotFacilityFilterService pilotFacilityFilterService;

    public StockRangeConsumer(StockRangeDataProcessor stockRangeDataProcessor,
                              @Qualifier(StockRangeConfiguration.STOCK_RANGE_CONSUMER) SqsMessageConsumer stockRangeConsumer,
                              ObjectMapper objectMapper,
                              PilotFacilityFilterService pilotFacilityFilterService) {
        this.pilotFacilityFilterService = pilotFacilityFilterService;
        this.objectMapper = objectMapper;
        this.stockRangeConsumer = stockRangeConsumer;
        this.stockRangeDataProcessor = stockRangeDataProcessor;
    }

    public void consume() throws ExecutionException, InterruptedException {
        List<Message> messages = stockRangeConsumer.retrieve();
        LOGGER.info("Number of available core in the processor is: {}", Runtime.getRuntime().availableProcessors());

        forkJoinPool.submit(() -> messages.parallelStream().forEach(message -> {
            try {
                RangePredictionData rangePredictionData = toRangePrediction(message.getMessage());
                LOGGER.info("Starting run analysis for part key: {}", rangePredictionData.getPartKey());
                LOGGER.debug("***** RangePredictionData with PartKey: {}, list size: {}", rangePredictionData.getPartKey(), messages.size());
                processRangePredictionData(rangePredictionData);
                LOGGER.info("Ended run analysis for part key: {}", rangePredictionData.getPartKey());
                stockRangeConsumer.delete(message);
            } catch (Exception e) {
                LOGGER.error("Error reading message from stock range queue : ", e);
            }
        })).get();

    }

    // TODO move logic to usecase/service and this class to adapter
    public void processRangePredictionData(RangePredictionData rangePredictionData) {
        if (rangePredictionData == null || rangePredictionData.getPartKey() == null || rangePredictionData.getRangeCalculationData() == null) {
            LOGGER.error("Range prediction consumed, but is null or missing range calculation data");
            return;
        }

        if (pilotFacilityFilterService.isPilotFacilityAndExcludedUsageCode(rangePredictionData.getPartKey().getFacilityCode(), rangePredictionData.getPartKey().getUsageCode())) {
            stockRangeDataProcessor.processRangePredictionData(rangePredictionData);
        }
    }

    private RangePredictionData toRangePrediction(String message) throws JsonProcessingException {
        return objectMapper.readValue(message, RangePredictionData.class);
    }

}

异步线程池配置

@Configuration
@EnableAsync
public class AsyncConfiguration {

    @Bean
    public Executor asyncExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(4);
        executor.setMaxPoolSize(10);
        return executor;
    }
}

调度线程池配置

@Configuration
@EnableScheduling
public class SchedulerConfiguration implements SchedulingConfigurer {
    

    @Override
    public void configureTasks(ScheduledTaskRegistrar taskRegistrar) {
        ThreadPoolTaskScheduler threadPoolTaskScheduler = new ThreadPoolTaskScheduler();

        threadPoolTaskScheduler.setPoolSize(8);
        threadPoolTaskScheduler.setThreadNamePrefix("scheduled-task-");
        threadPoolTaskScheduler.setRemoveOnCancelPolicy(true);
        threadPoolTaskScheduler.initialize();

        taskRegistrar.setTaskScheduler(threadPoolTaskScheduler);
    }
}

SQS客户端配置(waitTimeSeconds=0)

@Configuration
@Profile("aws")
public class StockRangeConfiguration {

    public static final String STOCK_RANGE_CONSUMER = "stockRangeSqsConsumer";

    private final String accessKey;
    private final String secretKey;
    private final String region;
    private final String queueUrl;
    private final int maxNumberOfMessages;
    private final String endpoint;

    public StockRangeConfiguration(@Value("${aws.sqs.stockRange.accessKey}") String accessKey,
                                   @Value("${aws.sqs.stockRange.secret_access_key}") String secretKey,
                                   @Value("${aws.sqs.stockRange.region}") String region,
                                   @Value("${aws.sqs.stockRange.queueUrl}") String queueUrl,
                                   @Value("${aws.sqs.stockRange.maxNumberOfMessages:10}") int maxNumberOfMessages,
                                   @Value("${aws.sqs.endpoint:}") String endpoint) {
        this.accessKey = accessKey;
        this.secretKey = secretKey;
        this.region = region;
        this.queueUrl = queueUrl;
        this.maxNumberOfMessages = maxNumberOfMessages;
        this.endpoint = endpoint;
    }

    @Bean(name = STOCK_RANGE_CONSUMER)
    public SqsMessageConsumer stockRangeConsumer() {
        if (endpoint.isEmpty()) {
            return new SqsMessageConsumer(new SqsConfiguration()
                    .withCredentials(accessKey, secretKey)
                    .withRegion(region)
                    .withUrl(queueUrl)
                    .withMaxNumberOfMessages(maxNumberOfMessages));
        }

        return new SqsMessageConsumer(new SqsConfiguration()
                .withCredentials(accessKey, secretKey)
                .withEndpoint(endpoint, region)
                .withUrl(queueUrl)
                .withMaxNumberOfMessages(maxNumberOfMessages));
    }
}

补充信息

调度线程池设为8,是因为存在其他调度任务,以100ms固定延迟从其他队列拉取10条消息。

优化建议

1. 启用SQS长轮询

当前waitTimeSeconds设为0属于短轮询,会频繁发起空请求浪费资源。将其设置为20秒(SQS最大长轮询时长),减少无效API调用,提升消息拉取效率,同时降低成本。

2. 消除调度线程阻塞

现有consume()方法中,forkJoinPool.submit(...).get()会阻塞调度线程,直到当前批次10条消息全部处理完成,导致无法及时发起下一轮拉取。改为异步非阻塞处理,让调度线程快速返回,立即触发下一次拉取:

// 注入全局异步线程池
private final Executor asyncExecutor;

// 构造函数补充注入
public StockRangeConsumer(..., Executor asyncExecutor) {
    this.asyncExecutor = asyncExecutor;
}

public void consume() {
    List<Message> messages = stockRangeConsumer.retrieve();
    if (messages.isEmpty()) {
        return;
    }
    messages.forEach(message -> asyncExecutor.submit(() -> {
        try {
            RangePredictionData rangePredictionData = toRangePrediction(message.getMessage());
            LOGGER.info("Starting run analysis for part key: {}", rangePredictionData.getPartKey());
            processRangePredictionData(rangePredictionData);
            LOGGER.info("Ended run analysis for part key: {}", rangePredictionData.getPartKey());
            stockRangeConsumer.delete(message);
        } catch (Exception e) {
            LOGGER.error("Error reading message from stock range queue : ", e);
        }
    }));
}

3. 优化线程池配置

4核ECS容器属于IO密集型场景(消息处理多涉及IO操作),调整线程池参数:

  • 异步线程池核心线程数设为8(4核*2),最大线程数设为16,避免线程过多导致上下文切换。
  • 调度线程池共8个线程,需确保当前消费者的调度任务有足够线程可用,可通过为该任务单独指定线程池,或调整其他任务的调度优先级。

4. 监控瓶颈环节

  • 监控ECS容器的CPU、内存、网络IO使用率,确认是否存在资源瓶颈。
  • 监控SQS队列ApproximateNumberOfMessagesVisible指标,跟踪消息堆积趋势。
  • 统计stockRangeDataProcessor.processRangePredictionData方法的耗时,若该方法是CPU密集型或依赖外部服务,需单独优化(如缓存结果、异步化调用、批量处理)。

5. 替换为Spring Cloud AWS原生监听

放弃自定义调度器,使用@SqsListener注解,原生支持并发消费、批量处理、自动重试等特性,比自定义实现更高效可靠:

@SqsListener(value = "${aws.sqs.stockRange.queueUrl}", maxNumberOfMessages = "10", waitTimeSeconds = "20")
public void listen(List<Message> messages) {
    messages.forEach(message -> {
        // 消息处理逻辑复用原有代码
    });
}

内容的提问来源于stack exchange,提问作者SamD

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 11:25:26