【问题标题】:@SQSListen results in an exception and not working@SQSListen 导致异常并且不起作用
【发布时间】:2020-03-15 06:16:47
【问题描述】:

我有一个非常简单的 Spring cloud aws 项目。我正在使用 Java 11。 这是 pom:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>2.2.5.RELEASE</version>
        <relativePath/> <!-- lookup parent from repository -->
    </parent>
    <groupId>com.demo.arf</groupId>
    <artifactId>testsqs-boot</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <name>testsqs-boot</name>
    <description>Demo project for Spring Boot</description>
    <dependencies>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-dependencies</artifactId>
            <version>Hoxton.SR3</version>
            <type>pom</type>
            <scope>runtime</scope>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-web</artifactId>
        </dependency>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-test</artifactId>
            <scope>test</scope>
            <exclusions>
                <exclusion>
                    <groupId>org.junit.vintage</groupId>
                    <artifactId>junit-vintage-engine</artifactId>
                </exclusion>
            </exclusions>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-aws</artifactId>
            <version>2.2.1.RELEASE</version>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-aws-messaging</artifactId>
            <version>2.2.1.RELEASE</version>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
            </plugin>
        </plugins>
    </build>
</project>

配置类:

package com.demo.arf.testsqsboot;

import com.amazonaws.services.sqs.AmazonSQSAsync;
import com.amazonaws.services.sqs.AmazonSQSAsyncClientBuilder;
import org.springframework.beans.factory.annotation.Value;
//import org.springframework.cloud.aws.messaging.core.QueueMessagingTemplate;
import org.springframework.cloud.aws.messaging.config.SimpleMessageListenerContainerFactory;
import org.springframework.cloud.aws.messaging.core.QueueMessagingTemplate;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import com.amazonaws.auth.AWSStaticCredentialsProvider;
import com.amazonaws.auth.BasicAWSCredentials;
import com.amazonaws.regions.Regions;
//import com.amazonaws.services.sqs.AmazonSQSAsync;
//import com.amazonaws.services.sqs.AmazonSQSAsyncClientBuilder;

@Configuration
public class SQSConfig {

    @Value("${cloud.aws.region.static}")
    private String region;

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

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

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

    public AmazonSQSAsync amazonSQSAsync() {
        return AmazonSQSAsyncClientBuilder.standard().withRegion(Regions.US_EAST_1)
                .withCredentials(new AWSStaticCredentialsProvider(new BasicAWSCredentials(awsAccessKey, awsSecretKey)))
                .build();
    }

    @Bean
    public SimpleMessageListenerContainerFactory simpleMessageListenerContainerFactory(AmazonSQSAsync amazonSQS){
        SimpleMessageListenerContainerFactory factory = new SimpleMessageListenerContainerFactory();
        factory.setAmazonSqs(amazonSQS);
        factory.setMaxNumberOfMessages(10);
        return factory;
    }
}

发送/接收消息的控制器类:

package com.demo.arf.testsqsboot.controller;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.aws.messaging.core.QueueMessagingTemplate;
import org.springframework.cloud.aws.messaging.listener.annotation.SqsListener;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class SQSController {
    @Autowired
    private QueueMessagingTemplate queueMessagingTemplate;

    private static final Logger LOG = LoggerFactory.getLogger(SQSController.class);
    @GetMapping("/send-sqs-message")
    public String sendMessage() {
        String sqsEndPoint= "https://sqs.us-east-2.amazonaws.com/1234567879/my_queue";
        queueMessagingTemplate.convertAndSend(sqsEndPoint, MessageBuilder.withPayload("hello from Spring Boot").build());
        return "Hello SQS";
    }

    @SqsListener("my_queue")
    public void getMessage(String message) {
      LOG.info(" *********** Message from SQS Queue - "+message);
    }
}

application.yml:

server:
  port: 9001
cloud:
  aws:
    region:
      static: us-east-1
      auto: false
    credentials:
      access-key: "asdmnasdn"
      secret-key: "sfkjsdjksdkj"
    end-point:
      uri: https://sqs.us-east-2.amazonaws.com/1234567879/my_queue

我可以让发送工作正常。但是当我添加监听器时,在启动过程中出现以下错误并且监听器没有收到消息:

2020-03-15 01:02:00.677  INFO 15423 --- [           main] o.s.web.context.ContextLoader            : Root WebApplicationContext: initialization completed in 3853 ms
2020-03-15 01:02:01.109  INFO 15423 --- [           main] o.s.s.concurrent.ThreadPoolTaskExecutor  : Initializing ExecutorService 'applicationTaskExecutor'
**WARNING: An illegal reflective access operation has occurred
WARNING: Illegal reflective access by com.amazonaws.util.XpathUtils (file:/Users/arf/.m2/repository/com/amazonaws/aws-java-sdk-core/1.11.415/aws-java-sdk-core-1.11.415.jar) to method com.sun.org.apache.xpath.internal.XPathContext.getDTMManager()
WARNING: Please consider reporting this to the maintainers of com.amazonaws.util.XpathUtils
WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations
WARNING: All illegal access operations will be denied in a future release
2020-03-15 01:02:01.749  WARN 15423 --- [           main]**
s.c.a.m.l.SimpleMessageListenerContainer : Ignoring queue with name 'my_queue': The queue does not exist.; nested exception is com.amazonaws.services.sqs.model.QueueDoesNotExistException: The specified queue does not exist for this wsdl version. (Service: AmazonSQS; Status Code: 400; Error Code: AWS.SimpleQueueService.NonExistentQueue; Request ID: 62821505-3f34-5434-a6ee)
2020-03-15 01:02:01.749  INFO 15423 --- [           main] o.s.s.concurrent.ThreadPoolTaskExecutor  : Initializing ExecutorService

还有一个基本问题。 @SQSListener 如何知道在哪里可以找到 aws 帐户信息和 sqs uri?

【问题讨论】:

    标签: java amazon-web-services spring-boot amazon-sqs


    【解决方案1】:

    我已经使它与配置类中的以下更改一起工作。但是,我想知道,大多数没有此代码的在线示例程序(使用 withEndpointConfiguration 构建 AmazonSQSAsync)是如何工作的。

        public QueueMessagingTemplate queueMessagingTemplate(AmazonSQSAsync amazonSQS) {
            return new QueueMessagingTemplate(amazonSQS);
        }
    
        @Bean
        @Primary
        public AmazonSQSAsync amazonSQS(AWSCredentialsProvider credentials) {
            return AmazonSQSAsyncClientBuilder.standard()
                    .withCredentials(credentials)
                    .withEndpointConfiguration(new AwsClientBuilder.EndpointConfiguration(localStackSqsUrl, awsRegion))
                    .build();
        }
    
        @Bean
        @Primary
        public AWSCredentialsProvider awsCredentialsProvider() {
            return new AWSCredentialsProviderChain(
                    new AWSStaticCredentialsProvider(
                            new BasicAWSCredentials("local", "stack")));
        }```
    

    【讨论】:

      【解决方案2】:

      几件事。
      - 首先永远不要将你的 AK/SK 存储在这样的属性或 yml 文件中。我可以看出这些是假值,但您总是希望从 ~/.aws/credentials 或实例元数据中提取这些值。如果您只需调用 .standard(),AWS 客户端(如 AmazonSqSAsyncClientBuilder)就会自动执行。无需凭据提供程序。
      - 其次,与地区相同
      - 第三,我相信@SqsListener 会使用你之前定义的ContainerFactory bean,至少@JmsListener 是这样工作的。
      - 您收到的错误消息是在您所选区域的帐户中未找到您的队列名称。您告诉它 us-east-1,但在您的发送代码中您指定了 us-east-2。根据您的帖子,我的猜测是您的队列在 us-east-2 中,因为您的问题是关于 @SqsListener,而不是 queueMessagingTemplate。

      【讨论】:

      • 我已经用 AWS 和 localstack 测试过这个。在这两种情况下,我都确保队列存在。但是我没有在本地设置 AWS 配置文件(~/.aws ...) 在侦听器上,SQSListener 和 JmsListener 的工作方式与我猜的相同。但是监听器怎么知道是哪个队列呢?特别是我有多个队列。在构建工厂时,我没有提供任何队列名称。
      • 但是您在 @SQSListener 注释中提供了队列名称...
      • 我有一些时间基本上复制了你正在做的事情。只要我在注释中使用队列 URL,我的工作就可以了
      • 队列URL在哪个注解中?我认为 @SQSListner 采用队列名称而不是完整的 URI。我的问题是 Listener。
      • Hi Sid - 队列 URL 位于 @SQSListener 注释上。文档说它将采用逻辑名称(来自 CloudFormation)、URL 或物理名称。我在github.com/kennyk65/spring-teaching-demos/tree/master/… 放了一些演示代码。看起来您的队列名称是正确的,但区域是错误的。
      猜你喜欢
      • 2012-06-13
      • 2020-03-21
      • 2013-09-28
      • 1970-01-01
      • 2015-08-24
      • 2014-09-06
      • 1970-01-01
      • 1970-01-01
      • 2010-12-15
      相关资源
      最近更新 更多