You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于WebSocket与ActiveMQ的Spring Boot推送通知系统技术问询

Hey Omar, great question—let’s tackle these two challenges one by one with practical, actionable solutions for your Angular + Spring Boot + ActiveMQ stack:

Challenge 1: Persist Notifications Until User Marks as Read

ActiveMQ’s default behavior removes messages once consumed, so we need to shift to manual message acknowledgment and per-user queue tracking to keep unread notifications available across logins. Here’s how to implement it:

Core Approach

  1. Use Manual Acknowledgment Mode: Configure your Spring Boot JMS listener to avoid auto-acknowledging messages. This keeps messages in the queue until explicitly confirmed (when the user marks them as read).
  2. Per-User Notification Queues: Create dedicated queues for each user (e.g., queue:user:{userId}) to isolate their unread notifications.
  3. Push Unread Notifications on Login: When a user logs in via Angular, your Spring Boot backend fetches all unacknowledged messages from their queue and pushes them over WebSocket.
  4. Acknowledge Only on "Mark as Read": When the user clicks to mark a notification as read, send a request to Spring Boot to trigger the message acknowledgment, removing it from the queue.

Code Examples

Spring Boot JMS Configuration (Manual ACK)

@Configuration
public class JmsConfig {
    @Bean
    public DefaultJmsListenerContainerFactory jmsListenerContainerFactory(ConnectionFactory connectionFactory) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        // Enable manual message acknowledgment
        factory.setSessionAcknowledgeMode(Session.CLIENT_ACKNOWLEDGE);
        return factory;
    }
}

Notification Service (Track & Acknowledge Messages)

@Service
public class NotificationService {
    // Track pending notifications per user (use a cache like Redis for distributed setups)
    private final Map<String, List<Message>> userPendingNotifications = new ConcurrentHashMap<>();

    @JmsListener(destination = "queue:user:${userId}", containerFactory = "jmsListenerContainerFactory")
    public void receiveNotification(Message message, Session session) throws JMSException {
        String userId = extractUserIdFromQueueName(message.getJMSDestination().toString());
        // Store message without acknowledging it (stays in queue)
        userPendingNotifications.computeIfAbsent(userId, k -> new ArrayList<>()).add(message);
    }

    // Push unread notifications when user connects via WebSocket
    public List<NotificationDto> getUnreadNotifications(String userId) {
        return userPendingNotifications.getOrDefault(userId, Collections.emptyList())
                .stream()
                .map(this::convertMessageToDto)
                .collect(Collectors.toList());
    }

    // Trigger acknowledgment when user marks as read
    public void markNotificationAsRead(String userId, String messageId) throws JMSException {
        List<Message> userMessages = userPendingNotifications.get(userId);
        if (userMessages == null) return;

        Optional<Message> targetMsg = userMessages.stream()
                .filter(msg -> msg.getJMSMessageID().equals(messageId))
                .findFirst();

        targetMsg.ifPresent(msg -> {
            try {
                msg.acknowledge(); // Remove message from queue
                userMessages.remove(msg);
            } catch (JMSException e) {
                // Handle acknowledgment failure (e.g., log and retry)
                e.printStackTrace();
            }
        });
    }

    private NotificationDto convertMessageToDto(Message message) throws JMSException {
        // Map JMS message to your DTO
        return new NotificationDto(
                message.getJMSMessageID(),
                message.getStringProperty("title"),
                message.getStringProperty("content")
        );
    }

    private String extractUserIdFromQueueName(String destinationName) {
        // Parse userId from queue name (e.g., "queue:user:123" → "123")
        return destinationName.split(":")[2];
    }
}

Angular WebSocket Service

@Injectable({ providedIn: 'root' })
export class NotificationWebSocketService {
  private socket: WebSocket | null = null;

  connect(userId: string): Observable<NotificationDto[]> {
    this.socket = new WebSocket(`ws://your-backend-url/notifications/${userId}`);
    
    return new Observable(observer => {
      this.socket?.onmessage = (event) => {
        const notifications = JSON.parse(event.data) as NotificationDto[];
        observer.next(notifications);
      };

      this.socket?.onerror = (error) => observer.error(error);
      this.socket?.onclose = () => observer.complete();
    });
  }

  markAsRead(messageId: string): void {
    this.http.post(`/api/notifications/mark-read/${messageId}`, {}).subscribe();
  }
}

Key Notes

  • Use ActiveMQ’s persistent storage (KahaDB or JDBC) to ensure notifications survive broker restarts.
  • Add message deduplication logic (via unique message IDs) to avoid re-pushing notifications if the client reconnects unexpectedly.

Challenge 2: Auto-Clean Idle Queues (1 Day Inactivity)

ActiveMQ doesn’t have built-in idle queue cleanup, but you can use its JMX API to automate this with a Spring Boot scheduled task.

Core Approach

  1. Schedule Daily Cleanup: Use Spring’s @Scheduled to run a cleanup job once per day.
  2. Query Queue Metadata via JMX: Connect to ActiveMQ’s JMX server to fetch queue statistics (last message received time, enqueue/dequeue counts).
  3. Delete Idle Queues: Remove any queue that hasn’t had message activity (send/receive) in the last 24 hours.

Code Example: Spring Boot Cleanup Service

@Service
public class IdleQueueCleanupService {
    private final MBeanServerConnection jmsMBeanServer;

    public IdleQueueCleanupService() throws IOException {
        // Connect to ActiveMQ's JMX server (default port: 1099)
        JMXServiceURL jmxUrl = new JMXServiceURL("service:jmx:rmi:///jndi/rmi://localhost:1099/jmxrmi");
        JMXConnector connector = JMXConnectorFactory.connect(jmxUrl);
        this.jmsMBeanServer = connector.getMBeanServerConnection();
    }

    @Scheduled(cron = "0 0 0 * * ?") // Run daily at midnight
    public void cleanupIdleQueues() throws Exception {
        long twentyFourHoursAgo = System.currentTimeMillis() - 24 * 60 * 60 * 1000;
        ObjectName queuePattern = new ObjectName("org.apache.activemq:type=Broker,brokerName=localhost,destinationType=Queue,destinationName=*");
        
        // Fetch all queue MBeans
        Set<ObjectName> queueNames = jmsMBeanServer.queryNames(queuePattern, null);

        for (ObjectName queueName : queueNames) {
            // Check last message activity time
            Long lastReceivedTime = (Long) jmsMBeanServer.getAttribute(queueName, "LastMessageReceivedTime");
            Long lastSentTime = (Long) jmsMBeanServer.getAttribute(queueName, "LastMessageSentTime");
            
            boolean isIdle = (lastReceivedTime == null || lastReceivedTime < twentyFourHoursAgo)
                    && (lastSentTime == null || lastSentTime < twentyFourHoursAgo);

            if (isIdle) {
                // Delete the idle queue
                jmsMBeanServer.invoke(queueName, "remove", new Object[]{}, new String[]{});
                System.out.println("Deleted idle queue: " + queueName.getKeyProperty("destinationName"));
            }
        }
    }
}

Key Notes

  • Ensure ActiveMQ has JMX enabled (default is true; verify in activemq.xml).
  • Add filters to avoid deleting system-critical queues (e.g., exclude queues not matching the queue:user:{userId} pattern).
  • Add logging/alerts to track deleted queues for debugging.

内容的提问来源于stack exchange,提问作者omark

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.15 04:57:40