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

求基于Spring Boot+Kafka+React的推送通知示例/参考方案

基于Kafka实现Spring Boot + React推送通知方案

一、整体流程设计

  • 业务事件触发时,Spring Boot后端将通知消息发送到Kafka的user-notifications Topic
  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 09:47:37