【问题标题】:ActiveMQ Concurrency Issue - Multiple Consumers Consuming the Same Message From QueueActiveMQ 并发问题 - 多个消费者从队列中消费相同的消息
【发布时间】:2015-02-02 01:15:46
【问题描述】:

我正在使用 Spring JMS 和 ActiveMQ,其中我有一个将消息推送到队列的客户端,并且我有多个消费者线程正在侦听并从队列中删除消息。有时 same 消息会被两个消费者从队列中出列。我不想要这种行为,并希望确保只有一个消息由一个消费者线程处理。关于我哪里出错的任何想法?

Spring 3.2.2 配置:

<beans xmlns="http://www.springframework.org/schema/beans"
    xmlns:context="http://www.springframework.org/schema/context"
    xmlns:mvc="http://www.springframework.org/schema/mvc" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    xmlns:tx="http://www.springframework.org/schema/tx" xmlns:util="http://www.springframework.org/schema/util"
    xsi:schemaLocation="
        http://www.springframework.org/schema/beans     
        http://www.springframework.org/schema/beans/spring-beans-3.0.xsd
        http://www.springframework.org/schema/context 
        http://www.springframework.org/schema/context/spring-context-3.0.xsd
        http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util-3.0.xsd
        http://www.springframework.org/schema/tx http://www.springframework.org/schema/tx/spring-tx-2.0.xsd">
    <context:annotation-config />
    <context:component-scan base-package="com.myapp" />

    <!-- JMS ConnectionFactory config Starts -->
    <bean id="jmsConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">
        <property name="brokerURL">
            <value>${brokerURL}</value>
        </property>
        <property name="userName" value="${username}" />
        <property name="password" value="${password}" />
    </bean>

    <bean id="pooledJmsConnectionFactory" class="org.apache.activemq.pool.PooledConnectionFactory"
        init-method="start" destroy-method="stop">
        <property name="connectionFactory" ref="jmsConnectionFactory" />
    </bean>
    <!-- JMS ConnectionFactory config Ends -->

    <!-- JMS Template config Starts -->
    <bean id="myQueue" class="org.apache.activemq.command.ActiveMQQueue">
        <constructor-arg value="${activemq.consumer.destinationName}" />
    </bean>

    <bean id="myQueueTemplate" class="org.springframework.jms.core.JmsTemplate">
        <property name="connectionFactory" ref="pooledJmsConnectionFactory" />
    </bean>
    <!-- JMS Template config Ends -->

    <!-- JMS Listener config starts -->
    <bean id="simpleMessageConverter"
        class="org.springframework.jms.support.converter.SimpleMessageConverter" />

    <bean id="myContainer" 
        class="org.springframework.jms.listener.DefaultMessageListenerContainer">
        <property name="concurrentConsumers" value="${threadcount}" />
        <property name="connectionFactory" ref="pooledJmsConnectionFactory" />
        <property name="destination" ref="myQueue" />
        <property name="messageListener" ref="myListener" />
        <property name="messageSelector" value="JMSType = 'New'" />
    </bean>

    <bean id="myListener"
        class="org.springframework.jms.listener.adapter.MessageListenerAdapter">
        <constructor-arg>
            <bean class="myapp.MessageListener" />
        </constructor-arg>
        <property name="defaultListenerMethod" value="receive" />
        <property name="messageConverter" ref="simpleMessageConverter" />
    </bean>
    <!-- JMS Listener config Ends -->


    <!-- enable the configuration of transactional behavior based on annotations -->
    <bean id="myJMSMessageSender" class="myapp.JMSMessageSender">
        <property name="jmsTemplate" ref="myQueueTemplate" />
        <property name="jmsQueue" ref="myQueue" />
        <property name="messageConverter" ref="simpleMessageConverter" />
    </bean>


    <bean id="myQueueTemplate" class="org.springframework.jms.core.JmsTemplate">
        <property name="connectionFactory" ref="pooledJmsConnectionFactory" />
    </bean>

</beans>

ActiveMQ 5.9.1 配置:

<broker xmlns="http://activemq.apache.org/schema/core" brokerName="instance8161" dataDirectory="${activemq.data}" persistent="false">

        <destinationPolicy>
            <policyMap>
              <policyEntries>
                <policyEntry topic="&gt;">
                    <!-- The constantPendingMessageLimitStrategy is used to prevent
                         slow topic consumers to block producers and affect other consumers
                         by limiting the number of messages that are retained
                         For more information, see:

                         http://activemq.apache.org/slow-consumer-handling.html

                    -->
                  <pendingMessageLimitStrategy>
                    <constantPendingMessageLimitStrategy limit="1000"/>
                  </pendingMessageLimitStrategy>
                </policyEntry>
              </policyEntries>
            </policyMap>
        </destinationPolicy>

        ... <!-- rest is default ActiveMQ Config -->
</broker>

【问题讨论】:

  • 您的消费者是否在确认发送两次的消息之前出错了?这当然可以解释消息的重新传递,如果是这样,您将不得不决定是多次获取消息(以保证您成功处理一次)还是不获取第二次(因此您永远不会成功)处理该消息)。

标签: java activemq spring-jms


【解决方案1】:

很可能,您的 myapp.MessageListener(或其依赖项之一)不是线程安全的,您会看到消费者线程之间的串扰。

最佳实践是将您的侦听器设计为无状态的(类中没有变异的字段)。如果这不可行,您需要使用锁来保护共享变量。

【讨论】:

  • 这些建议是为了避免可能导致此问题的事情之一,但消息侦听器中的线程安全问题并不是获得消息重新传递的唯一方法(也不是最有可能,根据我的经验),所以我会犹豫说这是“最有可能”的问题。这是一回事,但不要走这条路而排斥其他人。
  • 我同意;您关于重新传递的 cmets 是正确的,但是“有时相同的消息会被两个消费者从队列中出列。”在我看来,他是在谈论并发递送,而不是重新递送被拒绝的消息。
  • 这绝对是你解释它的方式。如果 OP 看到处理的消息总数是正确的,但有些被处理了不止一次,有些则从不处理,那么您的线程安全解释可能就是它。如果 OP 看到超过 N 次尝试处理 N 条消息,则可能是重新投递情况(或消息在某处重复的情况)。
猜你喜欢
  • 2012-08-11
  • 2021-11-25
  • 2016-06-08
  • 2014-01-22
  • 2013-08-01
  • 2013-06-09
  • 2014-06-27
  • 1970-01-01
  • 2023-03-28
相关资源
最近更新 更多