求基于Spring Boot+Kafka+React的推送通知示例/参考方案
基于Kafka实现Spring Boot + React推送通知方案
一、整体流程设计
- 业务事件触发时,Spring Boot后端将通知消息发送到Kafka的
user-notificationsTopic - Spring Boot推送服务作为Kafka消费者,监听该Topic获取通知消息
- 推送服务通过WebSocket将消息推送给在线的React前端用户
- React前端接收消息后展示通知
二、Spring Boot后端实现(生产+消费推送)
1. 依赖配置
在pom.xml中添加Kafka和WebSocket依赖:
<!-- Kafka --> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <!-- WebSocket --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency>
2. Kafka生产者配置(业务模块)
@Configuration public class KafkaProducerConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; @Bean public ProducerFactory<String, NotificationMessage> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, NotificationMessage> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } } // 通知消息实体 public class NotificationMessage implements Serializable { private String userId; private String content; private LocalDateTime timestamp; // getter/setter方法 } // 业务中发送通知 @Service public class NotificationService { @Autowired private KafkaTemplate<String, NotificationMessage> kafkaTemplate; public void sendUserNotification(String userId, String content) { NotificationMessage message = new NotificationMessage(); message.setUserId(userId); message.setContent(content); message.setTimestamp(LocalDateTime.now()); // 以userId为key,确保同一用户的消息被推送到同一个分区 kafkaTemplate.send("user-notifications", userId, message); } }
3. Kafka消费者+WebSocket推送服务
// WebSocket配置 @Configuration @EnableWebSocketMessageBroker public class WebSocketConfig implements WebSocketMessageBrokerConfigurer { @Override public void configureMessageBroker(MessageBrokerRegistry config) { // 启用简单消息代理,用于推送通知到客户端 config.enableSimpleBroker("/topic"); // 客户端发送消息的前缀 config.setApplicationDestinationPrefixes("/app"); } @Override public void registerStompEndpoints(StompEndpointRegistry registry) { // 注册WebSocket端点,客户端通过这个连接 registry.addEndpoint("/ws-notifications").setAllowedOrigins("*").withSockJS(); } } // Kafka消费者,消费消息后推送给WebSocket @Component public class NotificationConsumer { @Autowired private SimpMessagingTemplate messagingTemplate; @KafkaListener(topics = "user-notifications", groupId = "notification-push-group") public void consumeNotification(ConsumerRecord<String, NotificationMessage> record) { NotificationMessage message = record.value(); // 推送给指定用户的topic,比如/topic/notifications/{userId} messagingTemplate.convertAndSend("/topic/notifications/" + message.getUserId(), message); } }
三、React前端实现
1. 安装依赖
npm install sockjs-client stompjs
2. 通知订阅组件
import React, { useEffect, useState } from 'react'; import SockJS from 'sockjs-client'; import { Stomp } from 'stompjs'; const NotificationComponent = ({ userId }) => { const [notifications, setNotifications] = useState([]); let stompClient = null; useEffect(() => { // 建立WebSocket连接 const socket = new SockJS('http://localhost:8080/ws-notifications'); stompClient = Stomp.over(socket); stompClient.connect({}, () => { // 订阅当前用户的通知topic stompClient.subscribe(`/topic/notifications/${userId}`, (message) => { const notification = JSON.parse(message.body); setNotifications(prev => [...prev, notification]); }); }); return () => { if (stompClient) { stompClient.disconnect(); } }; }, [userId]); return ( <div className="notifications"> <h3>我的通知</h3> {notifications.map((item, index) => ( <div key={index} className="notification-item"> <p>{item.content}</p> <small>{new Date(item.timestamp).toLocaleString()}</small> </div> ))} </div> ); }; export default NotificationComponent;
四、关键注意事项
- 用户标识与消息路由:用
userId作为Kafka消息的key,确保同一用户的消息被消费后能精准推送到对应的前端连接 - 消息可靠性:配置Kafka的ACK机制(比如
acks=all)避免通知丢失;消费者可配置手动提交offset,确保消息处理完成后再确认 - 前端连接管理:实现WebSocket断开重连逻辑,结合用户登录状态维护连接有效性
- 消息序列化:使用JSON序列化消息,确保前后端能正确解析内容
参考思路扩展
- 低频次通知场景:前端定期轮询Spring Boot接口,接口内部从Kafka消费未读消息
- 大规模用户场景:引入Kafka Streams做消息聚合、过滤,降低推送服务压力
- 权限校验:结合Spring Security确保用户仅能订阅自身的通知topic
内容的提问来源于stack exchange,提问作者Dharmesh Vaghela
相关产品推荐
相关产品推荐

