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

如何正确配置多个DMLC监听单个SQS队列?性能未提升问题排查

订单管理系统SQS消息延迟问题排查

我们有一套订单管理系统,订单状态更新后会通过以下流程同步信息:先将消息发送至SQS队列,再由消费者调用客户端API。目前单条消息处理耗时约300-350ms,但SQS控制台显示最旧消息的延迟峰值可达50-60秒。

我判断单消费者无法承载当前负载,于是创建了多个DefaultMessageListenerContainer(DMLC)Bean,同时复制了多份消费者类,将这些消费者绑定到对应的DMLC中,但消息延迟情况没有改善,推测可能只有单个DMLC在处理消息,其余处于闲置状态。我是参考代码库中其他场景的做法添加多个DMLC的,不确定这是不是正确的解决方式。

消费者类代码

@Component
@Slf4j
@RequiredArgsConstructor
public class HOAEventsOMSConsumer extends ConsumerCommon implements MessageListener {

  private static final int MAX_RETRY_LIMIT = 3;
  private final OMSEventsWrapper omsEventsWrapper;
  @Override
  public void onMessage(Message message) {
    try {
      TextMessage textMessage = (TextMessage) message;
      String jmsMessageId = textMessage.getJMSMessageID();

      ConsumerLogging.logStart(jmsMessageId);

      String text = textMessage.getText();
      log.info(
          "Inside HOA Events consumer Request jmsMessageId:- " + jmsMessageId + " Text:- "
              + text);
      processAndAcknowledge(message, text, textMessage);
    } catch (JMSException e) {
      log.error("JMS Exception while processing surge message", e);
    }
  }

  private void processAndAcknowledge(Message message, String text, TextMessage textMessage) throws JMSException {
    try {
      TrimmedHOAEvent hoaEvent = JsonHelper.convertFromJsonPro(text, TrimmedHOAEvent.class);
      if (hoaEvent == null) {
        throw new OMSValidationException("Empty message in hoa events queue");
      }
      EventType event = EventType.fromString(textMessage.getStringProperty("eventType"));
      omsEventsWrapper.handleOmsEvent(event,hoaEvent);
      acknowledgeMessage(message);
    } catch (Exception e) {
      int retryCount = message.getIntProperty("JMSXDeliveryCount");
      log.info("Retrying... retryCount: {}, HOAEventsOMSConsumer: {}", retryCount, text);

      if (retryCount > MAX_RETRY_LIMIT) {
        log.info("about to acknowledge the message since it has exceeded maximum retry limit");
        acknowledgeMessage(message);
      }
    }
  }
}

DMLC配置类代码

@Configuration
@SuppressWarnings("unused")
public class HOAEventsOMSJMSConfig extends JMSConfigCommon{

  private Boolean isSQSQueueEnabled;

  @Autowired
  private HOAEventsOMSConsumer hoaEventsOMSConsumer;
  @Autowired
  private HOAEventsOMSConsumer2 hoaEventsOMSConsumer2;
  @Autowired
  private HOAEventsOMSConsumer3 hoaEventsOMSConsumer3;
  @Autowired
  private HOAEventsOMSConsumer4 hoaEventsOMSConsumer4;

  @Autowired
  private HOAEventsOMSConsumer5 hoaEventsOMSConsumer5;

  @Autowired
  private HOAEventsOMSConsumer6 hoaEventsOMSConsumer6;

  @Autowired
  private HOAEventsOMSConsumer7 hoaEventsOMSConsumer7;

  @Autowired
  private HOAEventsOMSConsumer8 hoaEventsOMSConsumer8;

  @Autowired
  private HOAEventsOMSConsumer9 hoaEventsOMSConsumer9;

  @Autowired
  private HOAEventsOMSConsumer10 hoaEventsOMSConsumer10;


  public HOAEventsOMSJMSConfig(IPropertyService propertyService, Environment env) {
    queueName = env.getProperty("aws.sqs.queue.oms.hoa.events.queue");
    endpoint = env.getProperty("aws.sqs.queue.endpoint") + queueName;
    JMSConfigCommon.accessId = env.getProperty("aws.sqs.access.id");
    JMSConfigCommon.accessKey = env.getProperty("aws.sqs.access.key");
    try {
      ServerNameCache serverNameCache = CacheManager.getInstance().getCache(ServerNameCache.class);
      if (serverNameCache == null) {
        serverNameCache = new ServerNameCache();
        serverNameCache.set(InetAddress.getLocalHost().getHostName());
        CacheManager.getInstance().setCache(serverNameCache);
      }
      this.isSQSQueueEnabled = propertyService.isConsumerEnabled(serverNameCache.get(), false);
    } catch (Exception e) {
      this.isSQSQueueEnabled = false;
    }
  }

  @Bean
  public JmsTemplate omsHOAEventsJMSTemplate(){
    SQSConnectionFactory sqsConnectionFactory;
    if (endpoint.toLowerCase().contains("localhost")) {
      sqsConnectionFactory =
          SQSConnectionFactory.builder().withEndpoint(getEndpoint("sqs")).build();
    } else {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withAWSCredentialsProvider(awsCredentialsProvider)
          .withNumberOfMessagesToPrefetch(10)
          .withEndpoint(endpoint)
          .build();
    }
    CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(sqsConnectionFactory);
    JmsTemplate jmsTemplate = new JmsTemplate(cachingConnectionFactory);
    jmsTemplate.setDefaultDestinationName(queueName);
    jmsTemplate.setDeliveryPersistent(false);
    jmsTemplate.setSessionTransacted(false);
    jmsTemplate.setSessionAcknowledgeMode(SQSSession.UNORDERED_ACKNOWLEDGE);
    return jmsTemplate;
  }

  @Bean
  public DefaultMessageListenerContainer jmsListenerHOAEventsListenerContainer() {
    SQSConnectionFactory sqsConnectionFactory;
    if (endpoint.toLowerCase().contains("localhost")) {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withEndpoint(getEndpoint("sqs"))
          .build();
    } else {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withAWSCredentialsProvider(awsCredentialsProvider)
          .withNumberOfMessagesToPrefetch(10)
          .withEndpoint(endpoint)
          .build();
    }
    DefaultMessageListenerContainer dmlc = new DefaultMessageListenerContainer();
    dmlc.setConnectionFactory(sqsConnectionFactory);
    dmlc.setDestinationName(queueName);
    dmlc.setAutoStartup(isSQSQueueEnabled);
    dmlc.setMessageListener(hoaEventsOMSConsumer);
    dmlc.setSessionTransacted(false);
    dmlc.setSessionAcknowledgeMode(SQSSession.UNORDERED_ACKNOWLEDGE);
    
    return dmlc;
  }

  @Bean
  public DefaultMessageListenerContainer jmsListenerHOAEventsListenerContainerNo2() {
    SQSConnectionFactory sqsConnectionFactory;
    if (endpoint.toLowerCase().contains("localhost")) {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withEndpoint(getEndpoint("sqs"))
          .build();
    } else {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withAWSCredentialsProvider(awsCredentialsProvider)
          .withNumberOfMessagesToPrefetch(10)
          .withEndpoint(endpoint)
          .build();
    }
    DefaultMessageListenerContainer dmlc = new DefaultMessageListenerContainer();
    dmlc.setConnectionFactory(sqsConnectionFactory);
    dmlc.setDestinationName(queueName);
    dmlc.setAutoStartup(isSQSQueueEnabled);
    dmlc.setMessageListener(hoaEventsOMSConsumer2);
    dmlc.setSessionTransacted(false);
    dmlc.setSessionAcknowledgeMode(SQSSession.UNORDERED_ACKNOWLEDGE);
    return dmlc;
  }

  @Bean
  public DefaultMessageListenerContainer jmsListenerHOAEventsListenerContainerNo3() {
    SQSConnectionFactory sqsConnectionFactory;
    if (endpoint.toLowerCase().contains("localhost")) {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withEndpoint(getEndpoint("sqs"))
          .build();
    } else {
      sqsConnectionFactory = SQSConnectionFactory.builder()
          .withAWSCredentialsProvider(awsCredentialsProvider)
          .withNumberOfMessagesToPrefetch(10)
          .withEndpoint(endpoint)
          .build();
    }
    DefaultMessageListenerContainer dmlc = new DefaultMessageListenerContainer();
    dmlc.setConnectionFactory(sqsConnectionFactory);
    dmlc.setDestinationName(queueName);
    dmlc.setAutoStartup(isSQSQueueEnabled);
    dmlc.setMessageListener(hoaEventsOMSConsumer3);
    dmlc.setSessionTransacted(false);
    dmlc.setSessionAcknowledgeMode(SQSSession.UNORDERED_ACKNOWLEDGE);
    return dmlc;
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 17:40:31