【问题标题】:Listen message queue SQS with Spring Boot not works with standard config使用 Spring Boot 侦听消息队列 SQS 不适用于标准配置
【发布时间】:2019-12-28 23:58:15
【问题描述】:

我无法使用 Spring Boot 和 SQS 创建工作队列侦听器 (消息发送并出现在 SQS ui 中)

@MessageMapping@SqsListener 不起作用

Java:11
Spring Boot:2.1.7
依赖:spring-cloud-aws-messaging

这是我的配置

@Configuration
@EnableSqs
public class SqsConfig {

    @Value("#{'${env.name:DEV}'}")
    private String envName;

    @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 Headers headers() {
        return new Headers();
    }

    @Bean
    public MessageQueue queueMessagingSqs(Headers headers,
                                          QueueMessagingTemplate queueMessagingTemplate) {
        Sqs queue = new Sqs();
        queue.setQueueMessagingTemplate(queueMessagingTemplate);
        queue.setHeaders(headers);
        return queue;
    }

    private ResourceIdResolver getResourceIdResolver() {
        return queueName -> envName + "-" + queueName;
    }

    @Bean
    public DestinationResolver destinationResolver(AmazonSQSAsync amazonSQSAsync) {
        DynamicQueueUrlDestinationResolver destinationResolver = new DynamicQueueUrlDestinationResolver(
                amazonSQSAsync,
                getResourceIdResolver());
        destinationResolver.setAutoCreate(true);
        return destinationResolver;
    }

    @Bean
    public QueueMessagingTemplate queueMessagingTemplate(AmazonSQSAsync amazonSQSAsync,
                                                         DestinationResolver destinationResolver) {
        return new QueueMessagingTemplate(amazonSQSAsync, destinationResolver, null);
    }

    @Bean
    public QueueMessageHandlerFactory queueMessageHandlerFactory() {
        QueueMessageHandlerFactory factory = new QueueMessageHandlerFactory();
        MappingJackson2MessageConverter messageConverter = new MappingJackson2MessageConverter();
        messageConverter.setStrictContentTypeMatch(false);
        factory.setArgumentResolvers(Collections.singletonList(new PayloadArgumentResolver(messageConverter)));
        return factory;
    }

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

}

我还注意到 org.springframework.cloud.aws.messaging.config.SimpleMessageListenerContainerFactoryorg.springframework.cloud.aws.messaging.config.annotation.SqsConfiguration 在启动时运行

还有我的测试

@RunWith(SpringJUnit4ClassRunner.class)
public class ListenTest {

    @Autowired
    private MessageQueue queue;

    private final String queueName = "test-queue-receive";

    private String result = null;

    @Test
    public void test_listen() {
        // given
        String data = "abc";

        // when
        queue.send(queueName, data).join();

        // then
        Awaitility.await()
                .atMost(10, TimeUnit.SECONDS)
                .until(() -> Objects.nonNull(result));

        Assertions.assertThat(result).equals(data);
    }

    @MessageMapping(value = queueName)
    public void receive(String data) {
        this.result = data;
    }
}

你觉得有什么不对吗?

我创建了一个例如 repo : (https://github.com/mmaryo/java-sqs-test)
在测试文件夹中,更改“application.yml”中的 aws 凭据
然后运行测试

【问题讨论】:

  • 请比“不起作用”更具体。具体会发生什么?是否有任何地方出现错误消息,或者 SQS 错误队列中有消息?
  • SQS 队列中的消息留在队列中,receive() 方法永远不会运行。好像@MessageMapping(value = queueName)不听队列?
  • 我不确定这个工具。我只用过@SqsListener
  • @SqsListener 也不起作用:/
  • 在 QueueMessageHandler SqsListener sqsListenerAnnotation = AnnotationUtils.findAnnotation(method, SqsListener.class); 中始终为空。所以 Spring dot 不扫描 @SpringBootTest 中的 @Sqs 也不扫描 @Compent

标签: java spring-boot amazon-sqs spring-messaging spring-cloud-aws


【解决方案1】:

我在使用 spring-cloud-aws-messaging 包时遇到了同样的问题,但后来我在 @SqsListener 注释中使用了队列 URL 而不是队列名称,并且它起作用了。

@SqsListener(value = { "https://full-queue-URL" }, deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS)
public void receive(String message) {
     // do something
}

在使用 spring-cloud-starter-aws-messaging 包时,您似乎可以使用队列名称。如果您不想使用入门包,我相信有一些配置允许使用队列名称而不是 URL。

编辑:尽管我在属性文件中列出了 us-east-1,但我注意到该区域默认为 us-west-2。然后我创建了一个 RegionProvider bean 并将区域设置为 us-east-1 ,现在当我在 @SqsMessaging 中使用队列名称时,它被找到并正确解析为框架代码中的 URL。

【讨论】:

  • 你能发布 RegionProvider @Bean 方法的 sn-p 吗?
  • @heug @Bean public RegionProvider regionProvider() { return () -> Region.getRegion(Regions.fromName(queueRegion)); }
【解决方案2】:

您需要利用 @Primary 注释,这对我有用:

@Autowired(required = false)
private AWSCredentialsProvider awsCredentialsProvider;

@Autowired
private AppConfig appConfig;

@Bean
public QueueMessagingTemplate getQueueMessagingTemplate() {
    return new QueueMessagingTemplate(sqsClient());
}

@Primary
@Bean
public AmazonSQSAsync sqsClient() {
    AmazonSQSAsyncClientBuilder builder = AmazonSQSAsyncClientBuilder.standard();

    if (this.awsCredentialsProvider != null) {
        builder.withCredentials(this.awsCredentialsProvider);
    }

    if (appConfig.getSqsRegion() != null) {
        builder.withRegion(appConfig.getSqsRegion());
    } else {
        builder.withRegion(Regions.DEFAULT_REGION);
    }

    return builder.build();
}

build.gradle 需要这些部门:

implementation("org.springframework.cloud:spring-cloud-starter-aws:2.2.0.RELEASE")
implementation("org.springframework.cloud:spring-cloud-aws-messaging:2.2.0.RELEASE")

【讨论】:

  • 什么是 AppConfig?
  • @EvanGertis 对包含环境变量的 bean 的引用
猜你喜欢
  • 2021-10-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-01-07
  • 2016-11-19
  • 2012-08-26
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多