【问题标题】:Spring rabbitmq message ordering not working anymoreSpring rabbitmq 消息排序不再起作用
【发布时间】:2021-06-27 14:00:34
【问题描述】:

我正在处理一个消息排序问题,在不久前修复了它之后,现在修复不再起作用了。

只是为了概述,我有以下环境:

顺序在 tcpAdapter 和消息接收者之间的某处丢失。

这个问题我已经解决了:

  1. 在生产者方面 - 使用发布者确认并返回
  rabbitmq:
    publisher-confirms: true
    publisher-returns: true
  1. 在消费者端 - 强制执行单线程执行器: 我在这里找到的想法:RabbitMQ - Message order of delivery,为此我使用了后处理器。
@Component
public class RabbitConnectionFactoryPostProcessor implements BeanPostProcessor {
  @Override
  public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
    if (bean instanceof CachingConnectionFactory) {
      ((CachingConnectionFactory) bean).setExecutor(Executors.newSingleThreadExecutor());
    }
    return bean;
  }
}

现在,在一些 master-pom 更新后(我们不控制 master pom,它处于项目级别),修复突然不再起作用。检查差异后,我没有看到 spring-rabbit 或 spring-amqp 有任何变化,我不明白为什么会有影响。


如果你想要具体的例子,这里有更多细节:

  1. 制片人。

TCP Server 向 tcpAdapter 应用程序发送消息,该应用程序使用 spring-integration 流从 TCP 获取消息并将其发送到 rabbitmq。

这是执行此操作的代码(inboundAdapterClient 我没有在此处发布,因为我认为它不重要):

  @Bean
  public IntegrationFlow tcpToRabbitFlowClient() {
    return IntegrationFlows.from(inboundAdapterClient())      
        .transform(tcpToRabbitTransformer)     
        .channel(TCP_ADAPTER_SOURCE);
        .get();
  }

tcpAdapter 应用程序以正确的顺序从 TCP 接收消息,但是 tcpAdapter rabbitmq 堆栈每次都不会以正确的顺序发送它们(80% 的时间正常,20% 的错误顺序)

这里是 spring boot yml 配置(仅相关信息):

spring:
  rabbitmq:
    publisher-confirms: true
    publisher-returns: true
  cloud:
    stream:
      bindings:
        tcpAdapterSource:
          binder: rabbit
          content-type: application/json
          destination: tcpadapter.messagereceiver
  1. 消费者。

消息接收者强制执行单线程执行器以及如下配置。

这里是spring boot yml配置(只有相关信息)

spring:
  cloud:     
    stream:
      bindings:
        fromTcpAdapter:
          binder: rabbit
          content-type: application/json
          destination: tcpadapter.messagereceiver
      rabbit:
        default:
          producer:
            exchangeDurable: false
            exchangeAutoDelete: true
          consumer:
            exchangeDurable: false
            exchangeAutoDelete: true

注意:只有一个生产者和一个消费者。

来自 pom 的一些版本,也许有帮助:

      <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot</artifactId>
        <version>2.2.4.RELEASE</version>
      </dependency>
      <dependency>
        <groupId>org.springframework.amqp</groupId>
        <artifactId>spring-amqp</artifactId>
        <version>2.2.3.RELEASE</version>
      </dependency>
      <dependency>
        <groupId>org.springframework.amqp</groupId>
        <artifactId>spring-rabbit</artifactId>
        <version>2.2.3.RELEASE</version>
      </dependency>
      <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-stream</artifactId>
        <version>3.0.1.RELEASE</version>
      </dependency>

【问题讨论】:

  • 您是否介意尝试不使用队列中的消息,而是确认它们确实按所需顺序放置?我的意思是,这可能不是您的消费者方面的事实,而是 Spring Cloud Stream 生产者方面的事实......
  • 之前的 spring-rabbit 版本是什么?使用发布者确认时,将通道返回到缓存现在会延迟,直到收到确认,这可能导致后续发送在不同的通道上进行,并且它们可能会乱序到达。如果你可以将 spring-rabbit 升级到 2.3.x,你可以使用ThreadChannelConnectionFactory,这样可以确保每个线程的所有发送始终使用相同的通道。
  • @ArtemBilan 所以我应该在我的应用程序和 rabbitmq 代理之间使用中间人来检查消息?如果我看到生产者的订单丢失,除了升级到 spring-rabbit 2.3.x 以使用 ThreadChannelConnectionFactory 之外,还有什么解决方案吗?
  • GaryRussell 的问题:您能告诉我如何使用 ThreadChannelConnectionFactory 吗?我可以在 application.yml 中为 spring 云流设置一个特殊属性来启用它吗?或者我应该通过@Configuration 类来做吗?您是否有任何示例,因为从文档中我看不到解决方案。非常感谢!
  • “中间人”只是一个 RabbitMQ 管理控制台:rabbitmq.com/management.html。在那里,您可以导航到队列并查看其内容。当然,如果你的消费者是 UP,消息将从队列中拉出,你不能做任何假设。

标签: spring-integration spring-cloud-stream spring-amqp


【解决方案1】:

通过删除 yml 配置并使用如下所述的显式 bean 声明和工厂配置来解决。唯一的问题是性能缓慢,但发布商确认这是预期的。

所以确实是生产者问题。

  @Bean
  public CachingConnectionFactory connectionFactory() {
    com.rabbitmq.client.ConnectionFactory connectionFactoryClient = new com.rabbitmq.client.ConnectionFactory();
    connectionFactoryClient.setUsername(username);
    connectionFactoryClient.setPassword(password);
    connectionFactoryClient.setHost(hostname);
    connectionFactoryClient.setVirtualHost(vhost);
    return new CachingConnectionFactory(connectionFactoryClient);
  }

  @Bean("rabbitTemplateAdapter")
  @Primary
  public RabbitTemplate rabbitTemplate(CachingConnectionFactory connectionFactory) {
    connectionFactory.setPublisherConfirmType(CORRELATED);
    connectionFactory.setPublisherReturns(true);
    RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
    rabbitTemplate.setMandatory(true);
    rabbitTemplate.setConfirmCallback((correlationData, ack, cause)
            -> log.debug("correlationData({}),ack({}),cause ({})", correlationData, ack, cause));
    rabbitTemplate.setReturnCallback((message, replyCode, replyText, exchange, routingKey)
            -> log.debug("exchange({}),route({}),replyCode({}),replyText({}),message:{}",
            exchange, routingKey, replyCode, replyText, message));
    return rabbitTemplate;
  }

对于发送消息:

rabbitTemplateAdapter.invoke(t -> {
      t.convertAndSend(
              exchange,
              DESTINATION,
              jsonMessage.getPayload(),
              m -> {outboundMapper().fromHeadersToRequest(jsonMessage.getHeaders(), m.getMessageProperties());
                return m;
              });
      t.waitForConfirmsOrDie(10_000);
      return true;
    });

我使用 spring rabbit 和 amqp 版本做到了这一点:

<dependency>
  <groupId>org.springframework.amqp</groupId>
  <artifactId>spring-rabbit</artifactId>
  <version>2.2.3.RELEASE</version>
</dependency>
<dependency>
  <groupId>org.springframework.amqp</groupId>
  <artifactId>spring-amqp</artifactId>
  <version>2.2.3.RELEASE</version>
</dependency>

Spring amqp 文档帮助很大,使用的技术称为“Scoped Operations”: https://docs.spring.io/spring-amqp/docs/2.2.7.RELEASE/reference/html/#scoped-operations

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-03-09
    • 1970-01-01
    • 2023-03-29
    • 1970-01-01
    • 2015-02-15
    • 2018-02-11
    • 2018-09-07
    相关资源
    最近更新 更多