【问题标题】:SpringCloud AWS - SQSListener annotated method not receiving messagesSpringCloud AWS - SQSListener 注释的方法不接收消息
【发布时间】:2022-01-28 20:55:38
【问题描述】:

我正在使用 Spring Cloud AWS 2.3.2 编写 SQS 发布者/消费者应用程序

<dependency>
      <groupId>io.awspring.cloud</groupId>
      <artifactId>spring-cloud-aws-messaging</artifactId>
      <version>2.3.2</version>
</dependency>

我已经到了可以成功将 msgs 发布到我的 SQS 的地步,但是我的 @SqsListener 注释方法不会消耗 msgs。我在这里查看了其他问答,但似乎都没有提供任何适当的见解来解决这个问题。

我在这里关注 API 文档:https://docs.awspring.io/spring-cloud-aws/docs/current/reference/html/index.html#annotation-driven-listener-endpoints

我的配置定义如下:

@Configuration
public class SqsMessagingConfig {

    @Value("${cloud.aws.credentials.secret-key}")
    private String secretKey;
    @Value("${cloud.aws.credentials.access-key}")
    private String accessKey;

    private AWSCredentialsProvider awsCredentialsProvider() {
        return new AWSStaticCredentialsProvider(new BasicAWSCredentials(accessKey,
                secretKey));
    }

    @Bean
    @Primary
    public AmazonSQSAsync amazonSQSAsync() {
        return AmazonSQSAsyncClientBuilder
                .standard()
                .withRegion("us-east-2")
                .withCredentials(awsCredentialsProvider())
                .build();
    }

    @Bean
    public QueueMessagingTemplate queueMessagingTemplate(AmazonSQSAsync amazonSQSAsync) {
        return new QueueMessagingTemplate(amazonSQSAsync);
    }

    @Bean
    public SimpleMessageListenerContainerFactory simpleMessageListenerContainerFactory(AmazonSQSAsync amazonSQSAsync) {
        SimpleMessageListenerContainerFactory factory = new SimpleMessageListenerContainerFactory();
        factory.setAmazonSqs(amazonSQSAsync);
        factory.setAutoStartup(true);
        factory.setMaxNumberOfMessages(10);

        return factory;
    }

 
    @Bean()
    public QueueMessageHandlerFactory queueMessageHandlerFactory(final ObjectMapper mapper, final AmazonSQSAsync amazonSQSAsync) {
        final QueueMessageHandlerFactory queueHandlerFactory = new QueueMessageHandlerFactory();
        queueHandlerFactory.setAmazonSqs(amazonSQSAsync);
        queueHandlerFactory.setArgumentResolvers(Collections.singletonList(new PayloadMethodArgumentResolver(jackson2MessageConverter(mapper))));
        return queueHandlerFactory;
    }

    private MessageConverter jackson2MessageConverter(final ObjectMapper mapper) {
        final MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
        converter.setObjectMapper(mapper);
        return converter;
    }
}

然后我的 SqsService 如下所示:

@Service
public class SqsQueueService {
    private static final Logger logger = LoggerFactory.getLogger(SqsQueueService.class);
    private final QueueMessagingTemplate queueMessagingTemplate;
    private final ObjectWriter objectWriter;
    private final String QUEUE_NAME = "SCHEDULES";

    public SqsQueueService(QueueMessagingTemplate queueMessagingTemplate, ObjectMapper mapper) {
        this.queueMessagingTemplate = queueMessagingTemplate;
        this.objectWriter = mapper.writer();
    }

    public void send(List<Schedule> schedules) {
        List<String> originatingIds = schedules.stream().map(Schedule::getOriginatingId).collect(Collectors.toList());
        try {
            Message<String> message = MessageBuilder.withPayload(objectWriter.writeValueAsString(schedules))
                    .build();

            this.queueMessagingTemplate.convertAndSend(QUEUE_NAME, message);
            logger.info("Successfully sent {} schedule(s) to SQS, with originatingId={}", schedules.size(),
                    originatingIds);
        } catch (Exception e) {
            logger.error("Failed to send the following schedule(s) to SQS=" + originatingIds, e);
        }

    }

    // NO_REDRIVE ensures we do not re-queue messages forever. They will be sent to a DLQ if they exceed maxReceiveCount
    @SqsListener(value = "SCHEDULES", deletionPolicy = SqsMessageDeletionPolicy.NO_REDRIVE)
    private void receiveMessage(List<Schedule> schedules) String partnerId) {
        List<String> originatingIds = schedules.stream().map(Schedule::getOriginatingId).collect(Collectors.toList());
        logger.info("Received request from SQS for originatingId={}", originatingIds);
        try {
            someService.createSchedules(schedules);
        } catch (Exception e) {
            throw new RuntimeException("An issue occurred during ingest for originatingId=" + originatingIds, e);
        }
    }

}

我还尝试了 aws-autoconfigured 依赖项,但这增加了很多额外的噪音,我仍然无法从 SQS 中使用它。希望有人能发现我在哪里搞砸/错过了什么。我直接从 spring 开发人员那里看到的文档表明我在做正确的事情,但显然情况并非如此。

将消息发送到队列后,我可以看到它正在等待被消费,但没有任何反应。非常感谢任何帮助。

【问题讨论】:

    标签: amazon-web-services amazon-sqs spring-cloud-aws sqslistener


    【解决方案1】:

    添加spring cloud aws autoconfigure依赖:

     <dependency>
          <groupId>io.awspring.cloud</groupId>
          <artifactId>spring-cloud-aws-autoconfigure</artifactId>
     </dependency>
    

    https://docs.awspring.io/spring-cloud-aws/docs/current/reference/html/index.html#maven-dependencies

    【讨论】:

      猜你喜欢
      • 2017-10-03
      • 2021-11-10
      • 1970-01-01
      • 2019-04-24
      • 1970-01-01
      • 2022-01-19
      • 2020-09-01
      • 1970-01-01
      • 2017-05-29
      相关资源
      最近更新 更多