基于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
- 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).
- Per-User Notification Queues: Create dedicated queues for each user (e.g.,
queue:user:{userId}) to isolate their unread notifications. - 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.
- 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
- Schedule Daily Cleanup: Use Spring’s
@Scheduledto run a cleanup job once per day. - Query Queue Metadata via JMX: Connect to ActiveMQ’s JMX server to fetch queue statistics (last message received time, enqueue/dequeue counts).
- 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
相关产品推荐
相关产品推荐

