【问题标题】:Spring AMQP AcknowledgeMode.AUTO works slowSpring AMQP AcknowledgeMode.AUTO 工作缓慢
【发布时间】:2017-12-27 16:31:02
【问题描述】:

我有一个生产者每秒向 RabbitMQ 发送 20 条消息,我还有一个消费者,它应该以与产生消息相同的速度接收消息。

我必须实现一些条件:

  1. 每秒产生和消费 20 条消息。
  2. 保存生产订单。
  3. 消息不应丢失(这就是我使用 AcknowledgeMode.AUTO 的原因)。

当我使用 Spring AMQP 实现(org.springframework.amqp.rabbit)时,我的消费者每秒最多处理 6 条消息。但是,如果我使用原生 AMQP 库 (com.rabbitmq.client),它每秒会使用 ack - auto 和 manual 处理所有 20 条消息。

问题是:

为什么消费者案例中的 Spring 实现工作如此缓慢,我该如何解决这个问题?

如果我设置 prefetchCount(20),它会根据需要工作,但我不能使用 prefetch,因为它会在拒绝情况下破坏订单。

春季amqp:

@Bean
public ConnectionFactory connectionFactory() {
    CachingConnectionFactory connectionFactory = new CachingConnectionFactory(rabbitMqServer);
    connectionFactory.setUsername(rabbitMqUsername);
    connectionFactory.setPassword(rabbitMqPassword);
    return connectionFactory;
}

...

private SimpleMessageListenerContainer createContainer(Queue queue, Receiver receiver, AcknowledgeMode acknowledgeMode) {
    SimpleMessageListenerContainer persistentListenerContainer = new SimpleMessageListenerContainer();
    persistentListenerContainer.setConnectionFactory(connectionFactory());
    persistentListenerContainer.setQueues(queue);
    persistentListenerContainer.setMessageListener(receiver);
    persistentListenerContainer.setAcknowledgeMode(AcknowledgeMode.AUTO);
    return persistentListenerContainer;
}

...

@Override
public void onMessage(Message message) {saveToDb}

【问题讨论】:

    标签: java spring amqp spring-amqp


    【解决方案1】:

    Spring AMQP(2.0 之前)默认 prefetch 为 1,正如您所说,即使在拒绝之后也能保证订单。

    本机客户端默认不应用basicQos(),这实际上意味着它具有无限预取。

    所以你不是在比较苹果和苹果。

    尝试使用原生客户端 channel.basicQos(1),您应该会看到与默认 spring amqp 设置类似的结果。

    编辑

    将苹果与苹果进行比较时,无论有无框架,我都会得到相似的结果...

    @SpringBootApplication
    public class So47995535Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So47995535Application.class, args).close();
        }
    
        private final CountDownLatch latch = new CountDownLatch(100);
    
        private int nativeCount;
    
        private int rlCount;
    
        @Bean
        public ApplicationRunner runner(ConnectionFactory factory, RabbitTemplate template,
                SimpleMessageListenerContainer container) {
            return args -> {
                for (int i = 0; i < 100; i++) {
                    template.convertAndSend("foo", "foo" + i);
                }
                container.start();
                Connection conn = factory.createConnection();
                Channel channel = conn.createChannel(false);
                channel.basicQos(1);
                channel.basicConsume("foo", new DefaultConsumer(channel) {
    
                    @Override
                    public void handleDelivery(String consumerTag, Envelope envelope, BasicProperties properties,
                            byte[] body) throws IOException {
                        System.out.println("native " + new String(body));
                        channel.basicAck(envelope.getDeliveryTag(), false);
                        nativeCount++;
                        latch.countDown();
                    }
    
                });
                latch.await(60, TimeUnit.SECONDS);
                System.out.println("Native: " + this.nativeCount + " LC: " + this.rlCount);
                channel.close();
                conn.close();
                container.stop();
            };
        }
    
        @Bean
        public SimpleMessageListenerContainer container(ConnectionFactory connectionFactory) {
            SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory);
            container.setQueueNames("foo");
            container.setPrefetchCount(1);
            container.setAutoStartup(false);
            container.setMessageListener((MessageListener) m -> {
                System.out.println("LC " + new String(m.getBody()));
                this.rlCount++;
                this.latch.countDown();
            });
            return container;
        }
    
    }
    

    和

    Native: 50 LC: 50
    

    【讨论】:

    • 当使用 basicQos(1) 时,无论哪种方式,我都会得到相同的结果 - 请参阅我的编辑。
    • 哦,我明白了。谢谢你的解释。据我了解,我的条件速度-可靠-顺序三个都不可能全部满足吧?
    • 除非您可以加快处理速度或加快网络速度,否则不会。通过我的测试(除了递增一个 int 并倒计时一个锁存器之外,在侦听器中根本不做任何工作),使用localhost 上的代理,即使 prefetch = 1,我也可以达到 > 4000 条消息/秒。所以我会建议你查看监听器代码的效率。
    猜你喜欢
    • 1970-01-01
    • 2017-02-25
    • 2013-10-01
    • 2018-11-21
    • 1970-01-01
    • 1970-01-01
    • 2020-09-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多