代码之家  ›  专栏  ›  技术社区  ›  Integrating Stuff

Spring-普通RabbitMQ比普通RabbitMQ+JMS快很多?

  •  0
  • Integrating Stuff  · 技术社区  · 7 年前

    我有两个Spring RabbitMq配置,一个使用RabbitTemplate,一个使用JmsTemplate。


    类AmqpMailIntegrationPerfTestConfig :

    @Configuration
    @ComponentScan(basePackages = {
        "com.test.perf.amqp.receiver",
        "com.test.perf.amqp.sender"
    })
    @EnableRabbit
    public class AmqpMailIntegrationPerfTestConfig {
    
        @Bean
        public DefaultClassMapper classMapper() {
            DefaultClassMapper classMapper = new DefaultClassMapper();
            Map<String, Class<?>> idClassMapping = new HashMap<>();
            idClassMapping.put("mail", MailMessage.class);
            classMapper.setIdClassMapping(idClassMapping);
            return classMapper;
        }
    
        @Bean
        public Jackson2JsonMessageConverter jsonMessageConverter() {
            Jackson2JsonMessageConverter jsonConverter = new Jackson2JsonMessageConverter();
            jsonConverter.setClassMapper(classMapper());
            return jsonConverter;
        }
    
        @Bean
        public RabbitTemplate myRabbitTemplate(ConnectionFactory connectionFactory) {
            final RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
            rabbitTemplate.setMessageConverter(jsonMessageConverter());
            return rabbitTemplate;
        }
    
        @Bean
        public ConnectionFactory createConnectionFactory(){
            CachingConnectionFactory connectionFactory = new CachingConnectionFactory("localhost");
            return connectionFactory;
        }
    
        @Bean
        Queue queue() {
            return new Queue(AmqpMailSenderImpl.QUEUE_NAME, false);
        }
    
        @Bean
        TopicExchange exchange() {
            return new TopicExchange(AmqpMailSenderImpl.TOPIC_EXCHANGE_NAME);
        }
    
        @Bean
        Binding binding(Queue queue, TopicExchange exchange) {
            return BindingBuilder.bind(queue).to(exchange).with(AmqpMailSenderImpl.ROUTING_KEY);
        }
    
        @Bean
        public AmqpAdmin amqpAdmin() {
            return new RabbitAdmin(createConnectionFactory());
        }
    
        @Bean
        public SimpleRabbitListenerContainerFactory myRabbitListenerContainerFactory() {
            SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
            factory.setConnectionFactory(createConnectionFactory());
            factory.setMaxConcurrentConsumers(5);
            factory.setMessageConverter(jsonMessageConverter());
            return factory;
        }
    
    }
    

    com.test.perf.amqp.sender包中的AMQPMailSenderPerImpl类

    @Component
    public class AmqpMailSenderPerfImpl implements MailSender {
    
        public static final String TOPIC_EXCHANGE_NAME = "mails-exchange";
        public static final String ROUTING_KEY = "mails";
    
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        @Override
        public boolean sendMail(MailMessage message) {
            rabbitTemplate.convertAndSend(TOPIC_EXCHANGE_NAME, ROUTING_KEY, message);
            return true;
        }
    }
    

    :

    @Component
    public class AmqpMailReceiverPerfImpl implements ReceivedDatesKeeper {
    
        private Logger logger = LoggerFactory.getLogger(getClass());
    
        private Map<String,Date> datesReceived = new HashMap<String, Date>();
    
        @RabbitListener(containerFactory = "myRabbitListenerContainerFactory", queues = AmqpMailSenderImpl.QUEUE_NAME)
        public void receiveMessage(MailMessage message) {
            logger.info("------ Received mail! ------\nmessage:" + message.getSubject());
            datesReceived.put(message.getSubject(), new Date());
        }
    
        public Map<String, Date> getDatesReceived() {
            return datesReceived;
        }
    
    }
    

    类JmsMailIntegrationPerfTestConfig :

    @Configuration
    @EnableJms
    @ComponentScan(basePackages = {
            "com.test.perf.jms.receiver",
            "com.test.jms.sender"
    })
    public class JmsMailIntegrationPerfTestConfig {
    
        @Bean
        public MessageConverter jacksonJmsMessageConverter() {
            MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter();
    
            Map<String,Class<?>> typeIdMappings = new HashMap<String,Class<?>>();
            typeIdMappings.put("mail", MailMessage.class);
            converter.setTypeIdMappings(typeIdMappings);
    
            converter.setTargetType(MessageType.TEXT);
            converter.setTypeIdPropertyName("_type");
    
            return converter;
        }
    
        @Bean
        public ConnectionFactory createConnectionFactory(){
            RMQConnectionFactory connectionFactory = new RMQConnectionFactory();
            connectionFactory.setUsername("guest");
            connectionFactory.setPassword("guest");
            connectionFactory.setVirtualHost("/");
            connectionFactory.setHost("localhost");
            connectionFactory.setPort(5672);
    
            return connectionFactory;
        }
    
        @Bean(name = "myJmsFactory")
        public JmsListenerContainerFactory<?> myFactory(ConnectionFactory connectionFactory) {
            DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
            factory.setConnectionFactory(connectionFactory);
            factory.setConcurrency("10-50");
            factory.setMessageConverter(jacksonJmsMessageConverter());
            return factory;
        }
    
        @Bean
        public Destination jmsDestination() {
            RMQDestination jmsDestination = new RMQDestination();
            jmsDestination.setDestinationName("myQueue");
            jmsDestination.setAmqp(false);
            jmsDestination.setAmqpQueueName("mails");
            return jmsDestination;
        }
    
        @Bean
        public JmsTemplate myJmsTemplate(ConnectionFactory connectionFactory) {
            final JmsTemplate jmsTemplate = new JmsTemplate(connectionFactory);
            jmsTemplate.setMessageConverter(jacksonJmsMessageConverter());
            return jmsTemplate;
        }
    
    }
    

    :

    @Component
    public class JmsMailSenderImpl implements MailSender {
    
        private Logger logger = LoggerFactory.getLogger(getClass());
    
        @Autowired
        private JmsTemplate jmsTemplate;
    
        @Override
        public boolean sendMail(MailMessage message) {
            logger.info("Sending message!");
            jmsTemplate.convertAndSend("mailbox", message);
    
            return false;
        }
    
    }
    

    包com.test.perf.jms.receiver中的JMSMAILReceivePerfImpl类

    @Component
    public class JmsMailReceiverPerfImpl implements ReceivedDatesKeeper {
    
        private Logger logger = LoggerFactory.getLogger(getClass());
    
        private Map<String,Date> datesReceived = new HashMap<String, Date>();
    
        @JmsListener(destination = "mailbox", containerFactory = "myJmsFactory", concurrency = "10")
        public void receiveMail(MailMessage message) {
            datesReceived.put(message.getSubject(), new Date());
            logger.info("Received <" + message.getSubject() + ">");
        }
    
        public Map<String, Date> getDatesReceived() {
            return datesReceived;
        }
    
    }
    

    我通过启动10个线程并让各个邮件发送者分别发送1000封邮件来测试上述配置。

    *所有消息的总吞吐量时间:3687ms *处理一条消息的时间:817ms

    *所有消息的总吞吐量时间:41653ms

    这似乎表明带有JmsTemplate的版本没有并行工作,或者至少没有以最佳方式使用资源。

    我们想要的是使用JmsTemplate获得与RabbitTemplate相同的吞吐量时间,因此我们可以使用JMS作为抽象层。

    1 回复  |  直到 7 年前
        1
  •  1
  •   Gary Russell    7 年前

    我能理解为什么消费者方面的速度较慢——这是因为 Consumer.receive() basicGet() 对于每条消息 @RabbitListener 容器使用 basicConsume 预取计数为250。

    在JMS发送端,您需要使用 CachingConnectionFactory

    尽管如此,速度还是慢了一点;我建议你问问rabbitmq用户Google组,rabbitmq工程师们在哪里闲逛。他们维护JMS客户机。