【问题标题】:Executing messages in same sequence after getting the messages from sqs fifo queue从 sqs fifo 队列中获取消息后以相同的顺序执行消息
【发布时间】:2018-10-06 06:11:35
【问题描述】:

我正在定义一个类似于下面定义的消息侦听器。消息是查询语句。消息按顺序推送到队列中,以便在推送到表时不会违反任何外键规则。当我将 maxnumberofmessages 设置为 10 时,我看到查询似乎是乱序执行的。但是,当我将其设置为 1 时,我看不到任何问题。当我将 maxnumberofmessages 设置为大于 1 的值时,如何确保消息以与队列中相同的顺序读取?

@Bean
public SimpleMessageListenerContainer simpleMessageListenerContainer(AmazonSQSAsync amazonSQSAsync) {
    SimpleMessageListenerContainer simpleMessageListenerContainer = new SimpleMessageListenerContainer();
    simpleMessageListenerContainer.setAmazonSqs);
    simpleMessageListenerContainer.setMessageHandler(queueMessageHandler());
    simpleMessageListenerContainer.setMaxNumberOfMessages(10);
    return simpleMessageListenerContainer;
}


@SqsListener(value = "${sqs.url}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS)
public void onMessage(String serviceData, @Header("MessageId") String messageId, @Header("ApproximateFirstReceiveTimestamp") String approximateFirstReceiveTimestamp) {
    repository.execute(serviceData);

}

修改后的代码

@Bean
public SimpleMessageListenerContainerFactory simpleMessageListenerContainerFactory() {
    SimpleMessageListenerContainerFactory msgListenerContainerFactory = new SimpleMessageListenerContainerFactory();
    msgListenerContainerFactory.setAmazonSqs(amazonSQSAsyncClient());
    msgListenerContainerFactory.setAutoStartup(false);
    msgListenerContainerFactory.setMaxNumberOfMessages(10);
    msgListenerContainerFactory.setTaskExecutor(threadPoolTaskExecutor());
    return msgListenerContainerFactory;
}


@Bean
public ThreadPoolTaskExecutor threadPoolTaskExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(1);
    executor.setMaxPoolSize(1);
    executor.setThreadNamePrefix("queueExecutor");
    executor.initialize();
    return executor;
}

@SqsListener(value = "${sqs.url}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS)
public void onMessage(String serviceData, @Header("MessageId") String messageId, @Header("ApproximateFirstReceiveTimestamp") String approximateFirstReceiveTimestamp) {
    repository.execute(serviceData);

}

【问题讨论】:

    标签: spring amazon-web-services spring-boot spring-cloud spring-messaging


    【解决方案1】:

    那是因为那里的逻辑是这样的:

       ReceiveMessageResult receiveMessageResult = getAmazonSqs().receiveMessage(this.queueAttributes.getReceiveMessageRequest());
       CountDownLatch messageBatchLatch = new CountDownLatch(receiveMessageResult.getMessages().size());
       for (Message message : receiveMessageResult.getMessages()) {
            if (isQueueRunning()) {
                 MessageExecutor messageExecutor = new MessageExecutor(this.logicalQueueName, message, this.queueAttributes);
                 getTaskExecutor().execute(new SignalExecutingRunnable(messageBatchLatch, messageExecutor));
            } else {
                messageBatchLatch.countDown();
            }
      }
    

    关注getTaskExecutor().execute()。这就是你的消息被转移到他们自己的线程执行的方式。

    您可以考虑将其重新配置为ThreadPoolTaskExecutor,池大小为1。并且您的所有消息都将在同一个线程上处理,因此将提供订单。

    【讨论】:

    • 谢谢@Artem Bilan。我发布的代码发生了什么?是创建多个线程吗?
    • 正确,默认的ThreadPoolTaskExecutor实际上是基于maxNumberOfMessages。所以,你有 10 个并行线程。
    • 太棒了。谢谢!
    • 我还有一个问题。 isQueueRunning() 是什么意思?
    • 我假设这些是 Spring 库,是否可以使用 Spring 实现相同的功能?
    猜你喜欢
    • 2019-11-07
    • 2018-08-29
    • 1970-01-01
    • 1970-01-01
    • 2019-11-07
    • 2014-10-10
    • 1970-01-01
    • 2021-05-01
    • 1970-01-01
    相关资源
    最近更新 更多