【问题标题】:DefaultJmsListenerContainerFactory and Concurrent Connections not shutting downDefaultJmsListenerContainerFactory 和并发连接未关闭
【发布时间】:2016-09-13 19:02:23
【问题描述】:

我正在使用 Spring 4.x 的 DefaultJmsListenerContainerFactory 连接到 ActiveMQ 队列,使用 @JmsListener 处理来自该队列的消息,然后将消息推送到同一 ActiveMQ 代理上的主题。

我为消费者/侦听器和生产者使用一个缓存连接工厂,并且我将缓存消费者设置为 false,以便我可以缓存生产者,但不能缓存消费者。我还将并发设置为 1-3,我希望在应用程序启动时队列中至少有 1 个消费者,并且随着消息的增加,消费者的数量将达到 3。消息减少了,我预计消费者的数量也会下降到 1。但是,如果我查看线程(defaultmessagelistenercontainer-2/3),它们处于等待状态,并且不会关闭。当负载下降时,预期的消费者数量也将关闭,这不是预期的行为吗?请在下面查看我的配置,如果此行为不是开箱即用的,请告诉我,如果我需要添加一些内容以使其正常工作,如我上面所述。

ApplicationContext.java

    @Bean
public DefaultJmsListenerContainerFactory jmsListenerContainerFactory() throws Throwable {
    DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();

    factory.setConnectionFactory(connectionFactory());
    factory.setConcurrency(environment.getProperty("jms.connections.concurrent"));
    factory.setSessionTransacted(environment.getProperty("jms.connections.transacted", Boolean.class));
    return factory;
}

@Bean
public CachingConnectionFactory connectionFactory(){
    RedeliveryPolicy redeliveryPolicy = new RedeliveryPolicy();
    redeliveryPolicy.setInitialRedeliveryDelay(environment.getProperty("jms.redelivery.initial-delay", Long.class));
    redeliveryPolicy.setRedeliveryDelay(environment.getProperty("jms.redelivery.delay", Long.class));
    redeliveryPolicy.setMaximumRedeliveries(environment.getProperty("jms.redelivery.maximum", Integer.class));
    redeliveryPolicy.setUseExponentialBackOff(environment.getProperty("jms.redelivery.use-exponential-back-off", Boolean.class));
    redeliveryPolicy.setBackOffMultiplier(environment.getProperty("jms.redelivery.back-off-multiplier", Double.class));

    ActiveMQConnectionFactory activeMQ = new ActiveMQConnectionFactory(environment.getProperty("jms.queue.username"), environment.getProperty("jms.queue.password"), environment.getProperty("jms.broker.endpoint"));
    activeMQ.setRedeliveryPolicy(redeliveryPolicy);
    activeMQ.setPrefetchPolicy(prefetchPolicy());

    CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(activeMQ);
    cachingConnectionFactory.setCacheConsumers(environment.getProperty("jms.connections.cache.consumers", Boolean.class));
    cachingConnectionFactory.setSessionCacheSize(environment.getProperty("jms.cache.size", Integer.class));
    return cachingConnectionFactory;
}

@Bean
public JmsMessagingTemplate jmsMessagingTemplate(){
    ActiveMQTopic activeMQ = new ActiveMQTopic(environment.getProperty("jms.queue.out"));

    JmsMessagingTemplate jmsMessagingTemplate = new JmsMessagingTemplate(connectionFactory());
    jmsMessagingTemplate.setDefaultDestination(activeMQ);

    return jmsMessagingTemplate;
}

application.properties

jms.connections.concurrent=1-3
jms.connections.prefetch=1000
jms.connections.transacted=true
jms.connections.cache.consumers=false
jms.redelivery.initial-delay=1000
jms.redelivery.delay=1000
jms.redelivery.maximum=5
jms.redelivery.use-exponential-back-off=true
jms.redelivery.back-off-multiplier=2
jms.cache.size=3
jms.queue.in=in.queue
jms.queue.out=out.queue
jms.broker.endpoint=failover:(tcp://localhost:61616)

【问题讨论】:

    标签: spring jms activemq spring-jms spring-messaging


    【解决方案1】:

    尝试设置maxMessagesPerTask > 0

    @Bean
    public DefaultJmsListenerContainerFactory jmsListenerContainerFactory() throws Throwable {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
    
        factory.setConnectionFactory(connectionFactory());
        factory.setMaxMessagesPerTask(1);
        factory.setConcurrency(environment.getProperty("jms.connections.concurrent"));
        factory.setSessionTransacted(environment.getProperty("jms.connections.transacted", Boolean.class));
        return factory;
    }
    

    您可以参考文档http://docs.spring.io/spring-framework/docs/4.3.x/javadoc-api/org/springframework/jms/listener/DefaultMessageListenerContainer.html#setMaxMessagesPerTask-int-

    jms.connections.prefetch=1000 表示如果您有 1000 条消息在 Q 上等待,那么您将只有 1 个线程开始处理这 1000 条消息。

    例如jms.connections.prefetch=1 意味着消息将被平等地分派给所有可用线程,但最好设置maxMessagesPerTask < 0,因为长期任务避免频繁的线程上下文切换。 http://activemq.apache.org/what-is-the-prefetch-limit-for.html

    【讨论】:

    • 谢谢!我试过了,它奏效了!但是,当我进行分析时,我注意到 dmlc 容器线程不断被重新创建,假设基于接收尝试和每个任务的最大消息的值。这只是一个观察,有人可以评论吗?
    • 是的,我给你举了一个例子,每个任务最多有 1 条消息,但 1 太低了,正如我提供的文档链接中所说,每个线程在 1 次尝试后都会死掉。例如,您需要将 max messagespertask 增加到 10 以允许 10 次尝试,并增加 IdleTaskExecutionLimit 以允许在达到最大尝试次数时重用线程。在文档docs.spring.io/spring-framework/docs/4.3.x/javadoc-api/org/… 中有很好的解释
    猜你喜欢
    • 1970-01-01
    • 2017-01-22
    • 1970-01-01
    • 1970-01-01
    • 2014-12-18
    • 2016-06-09
    • 2016-08-07
    • 1970-01-01
    • 2012-08-08
    相关资源
    最近更新 更多