【问题标题】:How can configure activemq ex. queue time to live , or keep the message in queue until period of time or action?如何配置activemq ex。排队等待时间,还是将消息保留在队列中,直到一段时间或采取行动?
【发布时间】:2018-07-02 07:27:29
【问题描述】:

我正在使用 spring boot 和 websocket 构建通知系统,我使用 ActiveMQ 为离线用户保留队列,它运行良好。

我需要编辑一些配置,例如生存时间,将消息保留在队列中直到用户阅读它,我不知道如何配置它?

下面是它的实现:

@Configuration
@EnableWebSocketMessageBroker 
public class WebSocketConfig extends AbstractWebSocketMessageBrokerConfigurer {

      @Override
        public void configureMessageBroker(MessageBrokerRegistry config) {
            /*config.enableSimpleBroker("/topic");
            config.setApplicationDestinationPrefixes("/app");*/

          config
            .setApplicationDestinationPrefixes("/app")
            .setUserDestinationPrefix("/user")
            .enableStompBrokerRelay("/topic","/queue","/user")

            .setRelayHost("localhost")
            .setRelayPort(61613)
            .setClientLogin("guest")
            .setClientPasscode("guest");


        }


        public void registerStompEndpoints(StompEndpointRegistry registry) {
            registry.addEndpoint("/websocket").withSockJS();
        }

}

还有:

@Service
public  class NotificationWebSocketService {
@Autowired
private SimpMessagingTemplate messagingTemplate;

public void initiateNotification(WebSocketNotification notificationData) throws InterruptedException {

messagingTemplate.convertAndSendToUser(notificationData.getUserID(), "/reply", notificationData.getMessage());

}
}

在调用 NotificationWebSocketService 后,它将在 activemq 中创建队列“/user/Johon/reply”,包含消息,当用户订阅此队列时将收到消息。

如何配置队列生存时间,将消息保留在队列中直到用户阅读它?

【问题讨论】:

  • 祝你好运@userdemo
  • “将消息保留在队列中,直到用户读取它”这是队列的原理。除非您使用主题,在这种情况下您需要使用队列...
  • @HassenBennour 是的,我这样做了,当用户订阅 stompClient.subscribe('/user/Johon/reply', function (greeting) { showGreeting(greeting.body); }); 收到的消息时,但我担心如何控制队列 '/user/Johon/reply' 的生存时间,以及是否有保留此消息
  • "stompClient.subscribe('/user/Johon/reply'' --> '/user/Johon/reply' 是主题而不是队列。这就是我不理解的原因您的担忧。
  • 您可以设置“过期”标头(以毫秒为单位)并使用带有标头支持的 convertAndSendToUser 版本。

标签: spring-boot activemq stomp spring-websocket server-push


【解决方案1】:

单元测试来说明如何设置用户队列中的消息过期。 需要tomcat-embedded、spring-messaging和active-mq

import org.apache.catalina.Context;
import org.apache.catalina.Wrapper;
import org.apache.catalina.connector.Connector;
import org.apache.catalina.startup.Tomcat;
import org.apache.coyote.http11.Http11NioProtocol;
import org.apache.tomcat.util.descriptor.web.ApplicationListener;
import org.apache.tomcat.websocket.server.WsContextListener;
import org.junit.AfterClass;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Configuration;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.simp.SimpMessagingTemplate;
import org.springframework.messaging.simp.config.ChannelRegistration;
import org.springframework.messaging.simp.config.MessageBrokerRegistry;
import org.springframework.messaging.simp.stomp.*;
import org.springframework.messaging.support.ChannelInterceptorAdapter;
import org.springframework.web.SpringServletContainerInitializer;
import org.springframework.web.WebApplicationInitializer;
import org.springframework.web.servlet.support.AbstractAnnotationConfigDispatcherServletInitializer;
import org.springframework.web.socket.WebSocketHttpHeaders;
import org.springframework.web.socket.client.standard.StandardWebSocketClient;
import org.springframework.web.socket.config.annotation.AbstractWebSocketMessageBrokerConfigurer;
import org.springframework.web.socket.config.annotation.EnableWebSocketMessageBroker;
import org.springframework.web.socket.config.annotation.StompEndpointRegistry;
import org.springframework.web.socket.messaging.WebSocketStompClient;
import org.springframework.web.socket.sockjs.client.SockJsClient;
import org.springframework.web.socket.sockjs.client.WebSocketTransport;

import java.io.File;
import java.io.IOException;
import java.lang.reflect.Type;
import java.util.*;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

import static java.util.concurrent.TimeUnit.SECONDS;

public class Test48402361 {

    private static final Logger logger = LoggerFactory.getLogger(Test48402361.class);

    private static TomcatWebSocketTestServer server = new TomcatWebSocketTestServer(33333);

    @BeforeClass
    public static void beforeClass() throws Exception {
        server.deployConfig(Config.class);
        server.start();
    }

    @AfterClass
    public static void afterClass() throws Exception {
        server.stop();
    }

    @Test
    public void testUser() throws Exception {

        WebSocketStompClient stompClient = new WebSocketStompClient(new SockJsClient(Collections.singletonList(new WebSocketTransport(new StandardWebSocketClient()))));

        BlockingQueue<String> blockingQueue = new LinkedBlockingQueue<>();
        StompSession session = stompClient
                .connect("ws://localhost:" + server.getPort() + "/test", new WebSocketHttpHeaders(), new StompSessionHandlerAdapter() {
                })
                .get();
        // waiting until message 2 expired
        Thread.sleep(3000);
        session.subscribe("/user/john/reply", new StompFrameHandler() {
            @Override
            public Type getPayloadType(StompHeaders headers) {
                return byte[].class;
            }

            @Override
            public void handleFrame(StompHeaders headers, Object payload) {
                String message = new String((byte[]) payload);
                logger.debug("message: {}, headers: {}", message, headers);
                blockingQueue.add(message);
            }
        });
        String message = blockingQueue.poll(1, SECONDS);
        Assert.assertEquals("1", message);
        message = blockingQueue.poll(1, SECONDS);
        Assert.assertEquals("3", message);

    }

    public static class Config extends AbstractAnnotationConfigDispatcherServletInitializer {

        @Override
        protected Class<?>[] getRootConfigClasses() {
            return new Class[] { };
        }

        @Override
        protected Class<?>[] getServletConfigClasses() {
            return new Class[] { Mvc.class };
        }

        @Override
        protected String[] getServletMappings() {
            return new String[] { "/" };
        }
    }

    @Configuration
    @EnableWebSocketMessageBroker
    public static class Mvc extends AbstractWebSocketMessageBrokerConfigurer {

        @Override
        public void registerStompEndpoints(StompEndpointRegistry stompEndpointRegistry) {

            stompEndpointRegistry.addEndpoint("/test")
                    .withSockJS()
                    .setWebSocketEnabled(true);
        }

        @Override
        public void configureMessageBroker(MessageBrokerRegistry registry) {
            registry.enableStompBrokerRelay("/user").setRelayHost("localhost").setRelayPort(61614);
        }

        @Autowired
        private SimpMessagingTemplate template;

        @Override
        public void configureClientInboundChannel(ChannelRegistration registration) {
            registration.setInterceptors(new ChannelInterceptorAdapter() {
                @Override
                public Message<?> preSend(Message<?> message, MessageChannel channel) {

                    StompHeaderAccessor sha = StompHeaderAccessor.wrap(message);
                    switch (sha.getCommand()) {
                        case CONNECT:
    // after connect we send 3 messages to user john, one will purged after 2 seconds.
                            template.convertAndSendToUser("john", "/reply", "1");
                            Map<String, Object> headers = new HashMap<>();
                            headers.put("expires", System.currentTimeMillis() + 2000);
                            template.convertAndSendToUser("john", "/reply", "2", headers);
                            template.convertAndSendToUser("john", "/reply", "3");
                            break;
                    }
                    return super.preSend(message, channel);
                }
            });
        }
    }

    public static class TomcatWebSocketTestServer {

        private static final ApplicationListener WS_APPLICATION_LISTENER =
                new ApplicationListener(WsContextListener.class.getName(), false);

        private final Tomcat tomcatServer;

        private final int port;

        private Context context;


        public TomcatWebSocketTestServer(int port) {

            this.port = port;

            Connector connector = new Connector(Http11NioProtocol.class.getName());
            connector.setPort(this.port);

            File baseDir = createTempDir("tomcat");
            String baseDirPath = baseDir.getAbsolutePath();

            this.tomcatServer = new Tomcat();
            this.tomcatServer.setBaseDir(baseDirPath);
            this.tomcatServer.setPort(this.port);
            this.tomcatServer.getService().addConnector(connector);
            this.tomcatServer.setConnector(connector);
        }

        private File createTempDir(String prefix) {
            try {
                File tempFolder = File.createTempFile(prefix + '.', "." + getPort());
                tempFolder.delete();
                tempFolder.mkdir();
                tempFolder.deleteOnExit();
                return tempFolder;
            } catch (IOException ex) {
                throw new RuntimeException("Unable to create temp directory", ex);
            }
        }

        public int getPort() {
            return this.port;
        }


        @SafeVarargs
        public final void deployConfig(Class<? extends WebApplicationInitializer>... initializers) {

            this.context = this.tomcatServer.addContext("", System.getProperty("java.io.tmpdir"));

            // Add Tomcat's DefaultServlet
            Wrapper defaultServlet = this.context.createWrapper();
            defaultServlet.setName("default");
            defaultServlet.setServletClass("org.apache.catalina.servlets.DefaultServlet");
            this.context.addChild(defaultServlet);

            // Ensure WebSocket support
            this.context.addApplicationListener(WS_APPLICATION_LISTENER);

            this.context.addServletContainerInitializer(
                    new SpringServletContainerInitializer(), new HashSet<>(Arrays.asList(initializers)));
        }

        public void start() throws Exception {
            this.tomcatServer.start();
        }

        public void stop() throws Exception {
            this.tomcatServer.stop();
        }

    }

}

【讨论】:

  • 我使用了 Map&lt;String, Object&gt; headers = new HashMap&lt;&gt;(); headers.put("expires", System.currentTimeMillis() + 2000); template.convertAndSendToUser("john", "/reply", "2", headers) 部分,2 秒后它工作完美,消息将被删除,我还有另一个问题,正如我上面提到的,我为 ex 构建 通知模块。当用户“john”登录时,他将订阅 '/user/Johon/reply' 队列中的所有消息,有一种方法可以使队列中的消息至少保留 10 条消息,并且如果有任何新消息标志来确定新消息和旧消息
  • 对不起,不明白您的要求。你想要一些历史数据吗? - 所以将读取的消息保存在数据库中,队列用于快速异步保证传递,并且没有 API 来更改消息状态并将其标记为已读。当您阅读并确认时,它就消失了。
【解决方案2】:

"stompClient.subscribe('/user/Johon/reply' --> '/user/Johon/reply' 是主题而不是队列。

如果您的 Stomp 客户端未连接到主题“/user/Johon/reply”,他将丢失发送到该主题的每条消息。

所以你的解决方案是:

  1. 将您的主题“/user/Johon/reply”转换为队列,以便消息无限期地保留在队列中或直到服务器结束处理消息。
  2. 使用追溯消费者和订阅恢复政策

追溯消费者只是一个普通的 JMS 主题消费者,他 表示在订阅开始时,每次尝试都应该是 用于及时返回并发送任何旧消息(或最后一条消息 发送该主题),消费者可能错过了。 http://activemq.apache.org/retroactive-consumer.html

订阅恢复政策允许您在以下情况下及时返回 您订阅了一个主题。 http://activemq.apache.org/subscription-recovery-policy.html

  1. 使用持久订阅者

长期离线的持久主题订阅者 在系统中通常不需要。这样做的原因是 代理需要保留所有发送到这些主题的消息 订户说。而这个消息堆积会随着时间的推移耗尽代理 例如,商店限制并导致整体放缓 系统。 http://activemq.apache.org/manage-durable-subscribers.html

具有 Stomp 的持久订阅者: http://activemq.apache.org/stomp.html#Stomp-ActiveMQExtensionstoSTOMP

CONNECT client-id string 指定 JMS clientID 用于 与 activemq.subcriptionName 组合以表示持久 订阅者。

关于TTL的一些解释

客户端可以为每个客户端指定一个以毫秒为单位的生存时间值 它发送的消息。该值定义了一个消息过期时间,即 消息的生存时间和发送时的 GMT 之和(对于 事务发送,这是客户端发送消息的时间,而不是 事务提交的时间)。

默认生存时间为 0,因此消息保留在队列中 无限期地或直到服务器端处理消息

更新

如果你想使用外部 ActiveMQ 代理

删除@EnableWebSocketMessageBroker 并添加到连接器下方的activemq.xml 并重新启动代理。

 <transportConnector name="stomp" uri="stomp://localhost:61613"/>

如果您想嵌入 ActiveMQ Broker,请将 bean 添加到您的 WebSocketConfig 中:

 @Bean(initMethod = "start", destroyMethod = "stop")
    public BrokerService broker() throws Exception {
        final BrokerService broker = new BrokerService();
        broker.addConnector("stomp://localhost:61613");    
        return broker;
    }

以及所需的依赖项

    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-activemq</artifactId>
    </dependency>
    <dependency>
        <groupId>org.apache.activemq</groupId>
        <artifactId>activemq-stomp</artifactId>
    </dependency>

完整示例 Spring Boot WebSocket with embedded ActiveMQ Broker

http://www.devglan.com/spring-boot/spring-boot-websocket-integration-example

【讨论】:

  • 您好,Hessen,谢谢您的回答,但是当连接到 websocket function connect() { var socket = new SockJS('/websocket'); stompClient = Stomp.over(socket); stompClient.connect({}, function (frame) { setConnected(true); console.log('Connected: ' + frame); stompClient.subscribe('/user/Johon/reply', function (greeting) { showGreeting(greeting.body); }); }); } 时有些奇怪的地方
  • 我正在寻找简单的实现来使用 activemq ,请您指教@Hassen Bennour
  • 查看我的更新。正如我从您的 jmx 屏幕截图中看到的那样,/user/Johon/reply 是一个队列,那么奇怪的是什么?你有没有向这个队列发送消息?
猜你喜欢
  • 2022-06-14
  • 2010-09-10
  • 1970-01-01
  • 2016-03-19
  • 2015-10-22
  • 2020-03-30
  • 2022-11-17
  • 2014-11-01
  • 2016-02-21
相关资源
最近更新 更多