【问题标题】:Persist message in ActiveMQ across server restart跨服务器重新启动在 ActiveMQ 中保留消息
【发布时间】:2016-01-12 14:03:21
【问题描述】:

我正在学习 Spring Integration JMS。由于 ActiveMQ 是一个消息代理。我指的是这里给出的项目-> http://www.javaworld.com/article/2142107/spring-framework/open-source-java-projects-spring-integration.html?page=2#

但我想知道如何在 ActiveMQ 中保留消息。我的意思是我启动了 ActiveMQ,然后使用 REST 客户端发送请求。我在 for 循环中调用 publishService.send( message ); 50 次,并且接收器端我有 10 秒的睡眠计时器。这样 50 条消息就会排队,并以 10 秒的间隔开始处理。

编辑:

看下面的截图:

它说 50 条消息已入队,其中 5 条已出队。

但是在这之间我停止了 ActiveMQ 服务器,到它消耗了 50 条消息中的 5 条消息,然后再次重新启动它。

但后来我期待它在 Messages Enqueued 列中显示剩余的 45 个。但是我可以在那里看到 0(见下面的屏幕截图),并且在服务器重新启动后它们都消失了,而无需保留剩余的 45 条消息。我该如何解决这个问题?

请看下面的配置:

<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
       xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
       xmlns:context="http://www.springframework.org/schema/context"
       xmlns:int="http://www.springframework.org/schema/integration"
       xmlns:int-jms="http://www.springframework.org/schema/integration/jms"
       xmlns:oxm="http://www.springframework.org/schema/oxm"
       xmlns:int-jme="http://www.springframework.org/schema/integration"
       xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-2.5.xsd
                http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context-2.5.xsd
                http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
                http://www.springframework.org/schema/integration/jms http://www.springframework.org/schema/integration/jms/spring-integration-jms.xsd
                http://www.springframework.org/schema/oxm http://www.springframework.org/schema/oxm/spring-oxm-3.0.xsd">


    <!-- Component scan to find all Spring components -->
    <context:component-scan base-package="com.geekcap.springintegrationexample" />

    <bean class="org.springframework.web.servlet.mvc.annotation.AnnotationMethodHandlerAdapter">
        <property name="order" value="1" />
        <property name="messageConverters">
            <list>
                <!-- Default converters -->
                <bean class="org.springframework.http.converter.StringHttpMessageConverter"/>
                <bean class="org.springframework.http.converter.FormHttpMessageConverter"/>
                <bean class="org.springframework.http.converter.ByteArrayHttpMessageConverter" />
                <bean class="org.springframework.http.converter.xml.SourceHttpMessageConverter"/>
                <bean class="org.springframework.http.converter.BufferedImageHttpMessageConverter"/>
                <bean class="org.springframework.http.converter.json.MappingJackson2HttpMessageConverter" />
            </list>
        </property>
    </bean>

    <!-- Define a channel to communicate out to a JMS Destination -->
    <int:channel id="topicChannel"/>

    <!-- Define the ActiveMQ connection factory -->
    <bean id="connectionFactory" class="org.apache.activemq.spring.ActiveMQConnectionFactory">
        <property name="brokerURL" value="tcp://localhost:61616"/>
    </bean>

    <!--
        Define an adaptor that route topicChannel messages to the myTopic topic; the outbound-channel-adapter
        automagically fines the configured connectionFactory bean (by naming convention
      -->
    <int-jms:outbound-channel-adapter channel="topicChannel"
                                      destination-name="topic.myTopic"
                                      pub-sub-domain="true" />

    <!-- Create a channel for a listener that will consume messages-->
    <int:channel id="listenerChannel" />

    <int-jms:message-driven-channel-adapter id="messageDrivenAdapter"
                                            channel="getPayloadChannel"
                                            destination-name="topic.myTopic"
                                            pub-sub-domain="true" />

    <int:service-activator input-channel="listenerChannel" ref="messageListenerImpl" method="processMessage" />

    <int:channel id="getPayloadChannel" />

    <int:service-activator input-channel="getPayloadChannel" output-channel="listenerChannel" ref="retrievePayloadServiceImpl" method="getPayload" />

</beans>

另请看代码:

我在 for 循环中一次发送消息的控制器:

@Controller
public class MessageController
{
    @Autowired
    private PublishService publishService;

    @RequestMapping( value = "/message", method = RequestMethod.POST )
    @ResponseBody
    public void postMessage( @RequestBody com.geekcap.springintegrationexample.model.Message message, HttpServletResponse response )
    {
        for(int i = 0; i < 50; i++){
            // Publish the message
            publishService.send( message );

            // Set the status to 201 because we created a new message
            response.setStatus( HttpStatus.CREATED.value() );
        }
    }

}

我已应用计时器的消费者代码:

@Service
public class MessageListenerImpl
{
    private static final Logger logger = Logger.getLogger( MessageListenerImpl.class );

    public void processMessage( String message )
    {
        try {
            Thread.sleep(10000);
            logger.info( "Received message: " + message );
            System.out.println( "MessageListener::::::Received message: " + message );
        } catch (InterruptedException e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }

    }
}

通过进一步搜索,我发现here 根据 JMS 规范,默认交付模式是持久的。但在我的情况下,它似乎不起作用。

请帮助我进行正确的配置,以便消息可以在代理失败时持续存在。

【问题讨论】:

  • @Erik Williams 你能帮我解决这个问题吗?
  • @Tim Bish 你能指出我正确的方向吗?

标签: jms activemq spring-integration message-queue spring-jms


【解决方案1】:

这通常不是问题,而是 activeMQ 的构建方式。

您可以在“ActiveMQ in action”一书中找到以下解释

  • 一旦消息被消息消费和确认 消费者,它通常会从代理的消息存储中删除。

因此,当您重新启动服务器时,它只会向您显示代理消息存储中的消息。在大多数情况下,您永远不需要查看已处理的消息。

希望这会有所帮助!

祝你好运!

【讨论】:

  • 感谢您的回复。你让我在概念上很清楚。其实我想错了方向。我不想保留已经消费的消息,而是希望 ActiveMQ 在服务器故障后保留尚未被消费者消费的剩余消息。我已经更新了问题。请看一下。
  • 你能看看你在 标签中的 'activemq.xml' 是否指定了类似 'persistent=false' 的东西?如果是,则将其更改为 pesistent=true
  • 不,我在 activemq.xml 中找不到 persistent=false
  • 这就是我的&lt;broker&gt; 标签看起来像&lt;broker xmlns="http://activemq.apache.org/schema/core" brokerName="localhost" dataDirectory="${activemq.data}"&gt;
  • 我应该添加pesistent=true 这样它看起来像&lt;broker xmlns="http://activemq.apache.org/schema/core" brokerName="localhost" dataDirectory="${activemq.data}" pesistent="true"&gt;
猜你喜欢
  • 2016-01-13
  • 2021-12-05
  • 2018-05-25
  • 2013-03-16
  • 2016-10-27
  • 2012-05-07
  • 2011-08-08
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多