ActiveMQ redelivery/死信隊列

在使用Message Queue的過程中,總會由于種種原因而導(dǎo)致消息失敗。一個經(jīng)典的場景是一個生成者向Queue中發(fā)消息,里面包含了一組郵件地址和郵件內(nèi)容。而消費者從Queue中將消息一條條讀出來,向指定郵件地址發(fā)送郵件。消費者在發(fā)送消息的過程中由于種種原因會導(dǎo)致失敗,比如網(wǎng)絡(luò)超時、當(dāng)前郵件服務(wù)器不可用等。這樣我們就希望建立一種機制,對于未發(fā)送成功的郵件再重新發(fā)送,也就是重新處理。重新處理超過一定次數(shù)還不成功,就放棄對該消息的處理,記錄下來,繼續(xù)對剩余消息進行處理。

ActiveMQ為我們實現(xiàn)了這一功能,叫做ReDelivery(重新投遞)。當(dāng)消費者在處理消息時有異常發(fā)生,會將消息重新放回Queue里,進行下一次處理。當(dāng)超過重試次數(shù)時,會給broker發(fā)送一個"Poison ack",這個消息被認(rèn)為是a poison pill(毒丸),這時broker會將這個消息發(fā)送到DLQ。

在以下四種情況中,ActiveMQ消息會被重發(fā)給客戶端/消費者:

  • 在一個事務(wù)session中,并且調(diào)用了session.rollback()方法。
  • 在一個事務(wù)session中,session.commit()之前調(diào)用了commit.close()。
  • 在session中使用CLIENT_ACKNOWLEDGE簽收模式,并且調(diào)用了session.recover()方法。
  • 在session中使用AUTO_ACKNOWLEDGE簽收模式,在異步(messageListener)消費消息情況下,如果onMessage方法異常且沒有被catch,此消息會被redelivery。

缺省情況下:持久消息過期,會被送到DLQ,非持久消息不會送到DLQ(不會redelivery)。

可以在connectionFactory中注入自定義的redeliveryPolicy來改變?nèi)笔?shù):

    <bean id="redeliveryPolicy" class="org.apache.activemq.RedeliveryPolicy">
        <!--是否在每次嘗試重新發(fā)送失敗后,增長這個等待時間-->
        <property name="useExponentialBackOff" value="true"></property>
        <!--重發(fā)次數(shù),默認(rèn)為6次-->
        <property name="maximumRedeliveries" value="5"></property>
        <!--重發(fā)時間間隔,默認(rèn)為1秒-->
        <property name="initialRedeliveryDelay" value="1000"></property>
        <!--第一次失敗后重新發(fā)送之前等待500毫秒,第二次失敗再等待500 * 2毫秒,這里的2就是value-->
        <property name="backOffMultiplier" value="2"></property>
        <!--最大傳送延遲,只在useExponentialBackOff為true時有效,當(dāng)重連間隔大于最大重連間隔時,以后每次重連間隔都為最大重連間隔。-->
        <property name="maximumRedeliveryDelay" value="1000"></property>
    </bean>


    <!-- 在ConnectionFactory中應(yīng)用這個Policy。 -->

    <bean id="targetConnectionFactory" class="org.apache.activemq.ActiveMQConnectionFactory">
        <property name="brokerURL" value="tcp://188.166.236.173:61616"/>         
        <property name="redeliveryPolicy" ref="activeMQRedeliveryPolicy"/>
        <!-- <property name="useAsyncSend" value="true"/> 默認(rèn)就是異步發(fā)送-->
    </bean>     

在ActiveMQ 服務(wù)端的conf/activemq.xmlzhong的broker節(jié)點下添加:

    <destinationPolicy>
        <policyMap>
          <policyEntries>
              <policyEntry queue=">">
                 <deadLetterStrategy>
                    <individualDeadLetterStrategy queuePrefix="DLQ." useQueueForQueueMessages="true" />
                    <!-- 如果不想將過期消息放到DLQ中
                    <sharedDeadLetterStrategy processExpired="false" />  
                    -->
                    <!-- 如果想將非持久消息放入DLQ
                    <sharedDeadLetterStrategy processNonPersistent="true" />
                    -->
                 </deadLetterStrategy>
              </policyEntry>

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

測試會重發(fā)消息(redelivery)的四種方法:

在一個事務(wù)session中,并且調(diào)用了session.rollback():

    <bean id="jmsQueueContainerForDLQ" class="org.springframework.jms.listener.DefaultMessageListenerContainer">
        <property name="connectionFactory" ref="connectionFactory" />
        <property name="destination" ref="queueDestination" />
        <property name="messageListener" ref="consumerMessageListenerDLQ" />
        <!-- 如果支持事物的話,在接收消息后rollback會重發(fā)消息,進入死信隊列,默認(rèn)為false -->
        <property name="sessionTransacted" value="true" />
    </bean>

    public class ConsumerMessageListenerDLQ implements SessionAwareMessageListener<TextMessage> {
        public void onMessage(TextMessage message, Session session) {
            if(message instanceof TextMessage) {
                TextMessage textMessage = (TextMessage) message;
                try {
                        String text = textMessage.getText();
                        System.out.println(String.format("Received: %s",text));
                        if ("i want to redelivery".equals(text)){
                            throw new JMSException("process failed to test redelivery and DLQ");
                    }

                } catch (JMSException e) {
                    System.out.println("there is JMS exception: " + e.getMessage() );
                    //throw JmsUtils.convertJmsAccessException(e);
                    try {
                        session.rollback();
                    } catch (JMSException e1) {
                        e1.printStackTrace();
                    }
                }
            }
        }
    }

在一個事務(wù)session中,session.commit()之前調(diào)用了commit.close():

    <bean id="jmsQueueContainerForDLQ" class="org.springframework.jms.listener.DefaultMessageListenerContainer">
        <property name="connectionFactory" ref="connectionFactory" />
        <property name="destination" ref="queueDestination" />
        <property name="messageListener" ref="consumerMessageListenerDLQ" />
        <!-- 如果支持事物的話,在接收消息后rollback會重發(fā)消息,進入死信隊列,默認(rèn)為false -->
        <property name="sessionTransacted" value="true" />
    </bean>

    public class ConsumerMessageListenerDLQ implements SessionAwareMessageListener<TextMessage> {
        public void onMessage(TextMessage message, Session session) {
            if(message instanceof TextMessage) {
                TextMessage textMessage = (TextMessage) message;
                try {
                        String text = textMessage.getText();
                        System.out.println(String.format("Received: %s",text));
                        if ("i want to redelivery".equals(text)){
                            throw new JMSException("process failed to test redelivery and DLQ");
                    }

                } catch (JMSException e) {
                    System.out.println("there is JMS exception: " + e.getMessage() );
                    //throw JmsUtils.convertJmsAccessException(e);
                    try {
                        session.close();
                    } catch (JMSException e1) {
                        e1.printStackTrace();
                    }
                }
            }
        }
    }

在session中使用CLIENT_ACKNOWLEDGE簽收模式,并且調(diào)用了session.recover():

    <bean id="jmsQueueContainerForDLQ" class="org.springframework.jms.listener.DefaultMessageListenerContainer">
        <property name="connectionFactory" ref="connectionFactory" />
        <property name="destination" ref="queueDestination" />
        <property name="messageListener" ref="consumerMessageListenerDLQ" />
        <!-- 自動應(yīng)答模式消息不會重發(fā),進入死信隊列 -->
        <property name="sessionAcknowledgeMode" value="2"/>
    </bean>

    public class ConsumerMessageListenerDLQ implements SessionAwareMessageListener<TextMessage> {
        public void onMessage(TextMessage message, Session session) {
            if(message instanceof TextMessage) {
                TextMessage textMessage = (TextMessage) message;
                try {
                        String text = textMessage.getText();
                        System.out.println(String.format("Received: %s",text));
                        if ("i want to redelivery".equals(text)){
                            throw new JMSException("process failed to test redelivery and DLQ");
                    }

                } catch (JMSException e) {
                    System.out.println("there is JMS exception: " + e.getMessage() );
                    //throw JmsUtils.convertJmsAccessException(e);
                    try {
                        session.recover();
                    } catch (JMSException e1) {
                        e1.printStackTrace();
                    }
                }
            }
        }
    }

在session中使用AUTO_ACKNOWLEDGE簽收模式,異步Listener的onMessage()異常未被捕捉:

    public class Listener implements MessageListener {
        public void onMessage(Message message) {
            int i = 8/0;//會導(dǎo)致redelivery
            try {
                if(message instanceof ActiveMQTextMessage){
                    ActiveMQTextMessage textMessage = (ActiveMQTextMessage) message;
                    System.out.println("收到的消息:" + textMessage.getText());                   }
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }

持久消息過期,會被送到DLQ:

    <!-- 定義JmsTemplate的Queue類型 -->
    <bean id="jmsQueueTemplate" class="org.springframework.jms.core.JmsTemplate">
        <!-- 這個connectionFactory對應(yīng)的是我們定義的Spring提供的那個ConnectionFactory對象 -->
        <constructor-arg ref="connectionFactory" />
        <property name="messageConverter" ref="messageConverter"></property>
        <!-- 非pub/sub模型(發(fā)布/訂閱),即隊列模式 -->
        <property name="pubSubDomain" value="false" />
        <!-- 發(fā)送模式  DeliveryMode.NON_PERSISTENT=1:非持久 ; DeliveryMode.PERSISTENT=2:持久-->
        <property name="deliveryMode" value="2" />
        <!-- 2秒后過期,這個對點對點模式有效 -->
        <property name="timeToLive" value="2000" />
    </bean>

Junit測試:

@Test
public void testDLQ() throws  Exception{
    //jpaUserService.findOne("3");
    for(int i = 1;i<=10;i++){
        queueProducer.sendMessages("我是第"+i+"個");
    }
    Thread.sleep(50000);
    System.out.print("全部執(zhí)行完畢!!!");
}
最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時請結(jié)合常識與多方信息審慎甄別。
平臺聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點,簡書系信息發(fā)布平臺,僅提供信息存儲服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

  • ActiveMQ 即時通訊服務(wù) 淺析http://www.cnblogs.com/hoojo/p/active_m...
    bboymonk閱讀 1,586評論 0 11
  • 消息中間件 消息中間件有很多的用途和優(yōu)點: 1. 將數(shù)據(jù)從一個應(yīng)用程序傳送到另一個應(yīng)用程序,或者從軟件的一個模塊傳...
    錯位的季節(jié)閱讀 882評論 0 1
  • Spring Cloud為開發(fā)人員提供了快速構(gòu)建分布式系統(tǒng)中一些常見模式的工具(例如配置管理,服務(wù)發(fā)現(xiàn),斷路器,智...
    卡卡羅2017閱讀 136,552評論 19 139
  • 用什么藥 讓我痊愈? 除非那個人 回心轉(zhuǎn)意 試著讓自己 死一回 依然喚不他一絲憐憫 死了又有何用? 在別人眼里 什...
    98ae0474329c閱讀 304評論 4 8

友情鏈接更多精彩內(nèi)容