最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

基于SpringBoot + Redis Pub/Sub實現(xiàn)跨實例SSE消息推送

 更新時間:2026年04月23日 09:22:55   作者:五阿哥永琪  
在構(gòu)建Web應(yīng)用時,消息推送是一個常見需求——比如站內(nèi)信、訂單狀態(tài)更新、告警通知等,SSE相比WebSocket更輕量,適合單向推送場景,本文介紹一種基于Redis Pub/Sub的解決方案,讓SSE連接能夠跨實例互通,需要的朋友可以參考下

引言

在構(gòu)建Web應(yīng)用時,消息推送是一個常見需求——比如站內(nèi)信、訂單狀態(tài)更新、告警通知等。SSE(Server-Sent Events)相比WebSocket更輕量,適合單向推送場景。但當服務(wù)部署多個實例時,問題就出現(xiàn)了:用戶A連接的是實例1,用戶B連接的是實例2,A發(fā)送的消息如何推送給B?本文介紹一種基于Redis Pub/Sub的解決方案,讓SSE連接能夠跨實例互通。

上圖展示了完整的消息流轉(zhuǎn)過程:

  1. 建立連接:用戶A和B分別連接到不同的后端實例,每個實例維護著自己的SSE連接池。
  2. 發(fā)送消息:用戶A發(fā)起私信請求,請求落在實例1上。
  3. Redis廣播:實例1將消息發(fā)布到Redis的station:message頻道,所有訂閱了該頻道的實例都會收到。
  4. 推送消息:實例2發(fā)現(xiàn)目標用戶B在自己身上,通過SSE連接將消息推送給B。實例1收到廣播后也會檢查,發(fā)現(xiàn)目標用戶不在自己身上,則直接忽略。

一、引入依賴

Spring Boot Web提供了SSE的支持(SseEmitter),而spring-boot-starter-data-redis則為我們帶來了Redis連接和Pub/Sub能力。commons-pool2是連接池,生產(chǎn)環(huán)境必備,避免頻繁創(chuàng)建連接帶來的性能損耗。

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
    <groupId>org.apache.commons</groupId>
    <artifactId>commons-pool2</artifactId>
</dependency>

二、SSE連接管理器

SseEmitterManager是整個方案的核心。這里有幾個設(shè)計點值得注意:

  • 支持多標簽頁:一個用戶可能打開多個瀏覽器標簽頁,所以用Map<String, Map<String, SseEmitter>>來組織,外層key是userId,內(nèi)層key是token(每個標簽頁唯一)。
  • 生命周期管理:通過onCompletion、onTimeout、onError回調(diào)來清理連接,防止內(nèi)存泄漏。
  • 不超時設(shè)置:new SseEmitter(0L)表示永不超時,你也可以根據(jù)業(yè)務(wù)需要設(shè)置一個合理的超時時間(如30分鐘)。
  • 全站廣播:broadcast()方法會遍歷所有在線用戶并推送,適合系統(tǒng)公告類消息。
@Component
@Slf4j
public class SseEmitterManager {
    /**
     * 用戶ID -> (連接token -> SseEmitter)
     * 一個用戶可能有多個瀏覽器標簽頁,用token區(qū)分
     */
    private final Map<String, Map<String, SseEmitter>> userEmitters = new ConcurrentHashMap<>();
    /**
     * 建立SSE連接
     * @param userId 用戶ID
     * @param token 連接標識(可用UUID)
     */
    public SseEmitter connect(String userId, String token) {
        // 超時時間設(shè)為0表示不超時,也可設(shè)置具體毫秒數(shù)
        SseEmitter emitter = new SseEmitter(0L);
        // 注冊回調(diào):連接關(guān)閉時清理
        emitter.onCompletion(() -> removeEmitter(userId, token));
        emitter.onTimeout(() -> removeEmitter(userId, token));
        emitter.onError(e -> removeEmitter(userId, token));
        // 存儲連接
        userEmitters.computeIfAbsent(userId, k -> new ConcurrentHashMap<>())
                    .put(token, emitter);
        log.info("SSE connected: userId={}, token={}, total users={}", 
                 userId, token, userEmitters.size());
        return emitter;
    }
    /**
     * 向指定用戶推送消息
     */
    public void sendToUser(String userId, String message) {
        Map<String, SseEmitter> emitters = userEmitters.get(userId);
        if (emitters == null || emitters.isEmpty()) {
            log.debug("User {} not online, message stored for later", userId);
            return;
        }
        // 向該用戶所有連接推送
        emitters.forEach((token, emitter) -> {
            try {
                emitter.send(SseEmitter.event()
                    .name("message")
                    .data(message));
            } catch (IOException e) {
                log.error("Send to user {} failed, removing emitter", userId, e);
                removeEmitter(userId, token);
            }
        });
    }
    /**
     * 全站廣播
     */
    public void broadcast(String message) {
        userEmitters.forEach((userId, emitters) -> {
            sendToUser(userId, message);
        });
        log.info("Broadcast message to {} users", userEmitters.size());
    }
    /**
     * 獲取當前在線人數(shù)
     */
    public int getOnlineCount() {
        return userEmitters.size();
    }
    private void removeEmitter(String userId, String token) {
        Map<String, SseEmitter> emitters = userEmitters.get(userId);
        if (emitters != null) {
            emitters.remove(token);
            if (emitters.isEmpty()) {
                userEmitters.remove(userId);
            }
        }
    }
}

三、Redis消息訂閱者

RedisMessageSubscriber實現(xiàn)了MessageListener接口,它會監(jiān)聽station:message頻道。當收到Redis消息時:

  • 將JSON反序列化為StationMessage對象。
  • 根據(jù)type字段判斷是私信還是廣播。
  • 調(diào)用SseEmitterManager的方法推送給目標用戶。

這里有個細節(jié):每個后端實例都會收到自己發(fā)布的消息,所以推送前需要判斷目標用戶是否在當前實例上——這個判斷邏輯其實隱含在sendToUser中:如果目標用戶不在本實例的連接池里,就直接返回,不會報錯。

@Component
@Slf4j
public class RedisMessageSubscriber implements MessageListener {
    
    @Autowired
    private SseEmitterManager sseEmitterManager;
    
    @Autowired
    private ObjectMapper objectMapper;
    
    @Override
    public void onMessage(Message message, byte[] pattern) {
        try {
            String channel = new String(message.getChannel());
            String body = new String(message.getBody());
            
            // 解析消息
            StationMessage msg = objectMapper.readValue(body, StationMessage.class);
            
            log.info("Received Redis message: channel={}, type={}, target={}", 
                     channel, msg.getType(), msg.getTargetUserId());
            
            // 根據(jù)消息類型分發(fā)
            if ("user".equals(msg.getType())) {
                // 私信:發(fā)給指定用戶
                sseEmitterManager.sendToUser(msg.getTargetUserId(), msg.getContent());
            } else if ("broadcast".equals(msg.getType())) {
                // 廣播:發(fā)給所有在線用戶
                sseEmitterManager.broadcast(msg.getContent());
            }
            
        } catch (Exception e) {
            log.error("Failed to process Redis message", e);
        }
    }
}

四、Redis配置

RedisMessageListenerContainer是Spring Data Redis提供的消息容器,它會自動管理訂閱和監(jiān)聽線程。注意這里訂閱的是station:message頻道,你可以根據(jù)業(yè)務(wù)需要定義多個頻道(比如station:notice、station:system)。

RedisTemplate的序列化配置也很重要:key使用StringRedisSerializer保證可讀性,value使用Jackson2JsonRedisSerializer來支持對象存儲。

@Configuration
public class RedisConfig {
    
    @Bean
    public RedisMessageListenerContainer redisMessageListenerContainer(
            RedisConnectionFactory connectionFactory,
            RedisMessageSubscriber subscriber) {
        
        RedisMessageListenerContainer container = new RedisMessageListenerContainer();
        container.setConnectionFactory(connectionFactory);
        // 訂閱站內(nèi)信頻道
        container.addMessageListener(subscriber, new ChannelTopic("station:message"));
        return container;
    }
    
    @Bean
    public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
        RedisTemplate<String, Object> template = new RedisTemplate<>();
        template.setConnectionFactory(factory);
        template.setKeySerializer(new StringRedisSerializer());
        template.setValueSerializer(new Jackson2JsonRedisSerializer<>(Object.class));
        return template;
    }
}

五、消息實體

StationMessage是消息的載體,在Redis中傳輸?shù)腏SON格式就對應(yīng)這個結(jié)構(gòu)。字段設(shè)計上:

  • id:消息唯一標識,可用于去重和消息歷史記錄。
  • type:區(qū)分私信和廣播,方便訂閱者做路由。
  • targetUserId:私信時的接收者,廣播時可為空。
  • timestamp:時間戳,客戶端可以用來做消息排序或展示發(fā)送時間。
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class StationMessage {
    private String id;           // 消息ID
    private String type;         // user:私信, broadcast:廣播
    private String targetUserId; // 私信時的目標用戶ID
    private String content;      // 消息內(nèi)容
    private String senderId;     // 發(fā)送者ID
    private Long timestamp;      // 時間戳
}

六、消息發(fā)送服務(wù)

MessageService封裝了發(fā)送邏輯。核心動作很簡單:構(gòu)造StationMessage -> 序列化為JSON -> redisTemplate.convertAndSend()。發(fā)布之后,所有實例的訂閱者都會收到消息,相當于Redis幫我們做了一個“廣播式”的跨實例通信。

這種設(shè)計的優(yōu)點是:發(fā)送方不需要知道消息最終由哪個實例處理,也不需要維護實例之間的網(wǎng)絡(luò)連接,所有協(xié)調(diào)工作都交給了Redis。

@Service
@Slf4j
public class MessageService {
    
    @Autowired
    private RedisTemplate<String, Object> redisTemplate;
    
    @Autowired
    private SseEmitterManager sseEmitterManager;
    
    @Autowired
    private ObjectMapper objectMapper;
    
    /**
     * 發(fā)送私信
     */
    public void sendPrivateMessage(String fromUserId, String toUserId, String content) {
        StationMessage msg = StationMessage.builder()
            .id(UUID.randomUUID().toString())
            .type("user")
            .targetUserId(toUserId)
            .senderId(fromUserId)
            .content(content)
            .timestamp(System.currentTimeMillis())
            .build();
        
        try {
            String json = objectMapper.writeValueAsString(msg);
            // 發(fā)布到Redis,所有實例都會收到
            redisTemplate.convertAndSend("station:message", json);
            log.info("Private message sent: {} -> {}", fromUserId, toUserId);
        } catch (JsonProcessingException e) {
            log.error("Failed to serialize message", e);
        }
    }
    
    /**
     * 全站廣播
     */
    public void broadcast(String fromUserId, String content) {
        StationMessage msg = StationMessage.builder()
            .id(UUID.randomUUID().toString())
            .type("broadcast")
            .senderId(fromUserId)
            .content(content)
            .timestamp(System.currentTimeMillis())
            .build();
        
        try {
            String json = objectMapper.writeValueAsString(msg);
            redisTemplate.convertAndSend("station:message", json);
            log.info("Broadcast message sent by: {}", fromUserId);
        } catch (JsonProcessingException e) {
            log.error("Failed to serialize broadcast", e);
        }
    }
}

七、Controller層

Controller對外暴露了三個核心接口:

  • GET /api/sse/connect:客戶端通過EventSource或fetch API調(diào)用這個接口建立SSE連接。每個連接會生成一個唯一token,用于后續(xù)的清理。
  • POST /api/sse/private:發(fā)送私信,需要提供發(fā)送者、接收者和消息內(nèi)容。
  • POST /api/sse/broadcast:發(fā)送廣播消息。
  • GET /api/sse/online-count:查詢當前在線人數(shù),可用于展示“在線狀態(tài)”或做監(jiān)控。

生產(chǎn)環(huán)境建議給這些接口加上認證鑒權(quán)(比如從token中解析userId),避免偽造身份。

@RestController
@RequestMapping("/api/sse")
@Slf4j
public class SseController {
    
    @Autowired
    private SseEmitterManager sseEmitterManager;
    
    @Autowired
    private MessageService messageService;
    
    /**
     * SSE連接端點
     */
    @GetMapping(value = "/connect", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public SseEmitter connect(@RequestParam String userId) {
        String token = UUID.randomUUID().toString();
        return sseEmitterManager.connect(userId, token);
    }
    
    /**
     * 發(fā)送私信
     */
    @PostMapping("/private")
    public ResponseEntity<?> sendPrivate(@RequestBody PrivateMessageRequest request) {
        messageService.sendPrivateMessage(
            request.getFromUserId(),
            request.getToUserId(),
            request.getContent()
        );
        return ResponseEntity.ok().build();
    }
    
    /**
     * 全站廣播
     */
    @PostMapping("/broadcast")
    public ResponseEntity<?> broadcast(@RequestBody BroadcastRequest request) {
        messageService.broadcast(request.getFromUserId(), request.getContent());
        return ResponseEntity.ok().build();
    }
    
    /**
     * 獲取在線人數(shù)
     */
    @GetMapping("/online-count")
    public ResponseEntity<Integer> getOnlineCount() {
        return ResponseEntity.ok(sseEmitterManager.getOnlineCount());
    }
}

以上就是基于SpringBoot + Redis Pub/Sub實現(xiàn)跨實例SSE消息推送的詳細內(nèi)容,更多關(guān)于SpringBoot跨實例SSE消息推送的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Java使用fastjson對String、JSONObject、JSONArray相互轉(zhuǎn)換

    Java使用fastjson對String、JSONObject、JSONArray相互轉(zhuǎn)換

    這篇文章主要介紹了Java使用fastjson對String、JSONObject、JSONArray相互轉(zhuǎn)換,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-11-11
  • 淺談JVM中的JOL

    淺談JVM中的JOL

    我們天天都在使用java來new對象,但估計很少有人知道new出來的對象到底長的什么樣子?對于普通的java程序員來說,可能從來沒有考慮過java中對象的問題,不懂這些也可以寫好代碼。今天,給大家介紹一款工具JOL,可以滿足大家對java對象的所有想象。
    2021-06-06
  • Java?String類和StringBuffer類的區(qū)別介紹

    Java?String類和StringBuffer類的區(qū)別介紹

    這篇文章主要介紹了Java?String類和StringBuffer類的區(qū)別,?關(guān)于java的字符串處理我們一般使用String類和StringBuffer類有什么不同呢,下面我們一起來看看詳細介紹吧
    2022-03-03
  • Spring Cloud 專題之Sleuth 服務(wù)跟蹤實現(xiàn)方法

    Spring Cloud 專題之Sleuth 服務(wù)跟蹤實現(xiàn)方法

    這篇文章主要介紹了Spring Cloud 專題之Sleuth 服務(wù)跟蹤,本文通過實例代碼給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2021-08-08
  • Java中的Kotlin?內(nèi)部類原理

    Java中的Kotlin?內(nèi)部類原理

    這篇文章主要介紹了Java中的Kotlin?內(nèi)部類原理,文章圍繞主題展開詳細的內(nèi)容介紹,具有一定的參考價值,感興趣的小伙伴可以參考一下
    2022-06-06
  • JAVA--HashMap熱門面試題

    JAVA--HashMap熱門面試題

    這篇文章主要介紹了JAVA關(guān)于HashMap容易被提問的面試題,文中題目提問頻率高,相信對你的面試有一定幫助,想要入職JAVA的朋友可以了解下
    2020-06-06
  • JDK動態(tài)代理接口和接口實現(xiàn)類深入詳解

    JDK動態(tài)代理接口和接口實現(xiàn)類深入詳解

    這篇文章主要介紹了JDK動態(tài)代理接口和接口實現(xiàn)類,JDK動態(tài)代理是代理模式的一種實現(xiàn)方式,因為它是基于接口來做代理的,所以也常被稱為接口代理,文中通過實例代碼介紹的非常詳細,需要的朋友可以參考下
    2022-06-06
  • 淺談@PostConstruct不被調(diào)用的原因

    淺談@PostConstruct不被調(diào)用的原因

    這篇文章主要介紹了淺談@PostConstruct不被調(diào)用的原因及分析,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • 關(guān)于Java 項目封裝sqlite連接池操作持久化數(shù)據(jù)的方法

    關(guān)于Java 項目封裝sqlite連接池操作持久化數(shù)據(jù)的方法

    這篇文章主要介紹了Java 項目封裝sqlite連接池操作持久化數(shù)據(jù)的方法,文中給大家介紹了sqlite的體系結(jié)構(gòu)及封裝java的sqlite連接池的詳細過程,需要的朋友可以參考下
    2021-11-11
  • Java BufferWriter寫文件寫不進去或缺失數(shù)據(jù)的解決

    Java BufferWriter寫文件寫不進去或缺失數(shù)據(jù)的解決

    這篇文章主要介紹了Java BufferWriter寫文件寫不進去或缺失數(shù)據(jù)的解決方案,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-07-07

最新評論

邯郸市| 忻城县| 柳河县| 容城县| 万载县| 罗定市| 宣化县| 会宁县| 寻乌县| 姜堰市| 吴桥县| 昌宁县| 阳朔县| 长子县| 武隆县| 凤山市| 仲巴县| 嵊州市| 洛川县| 兴安盟| 博罗县| 阜康市| 肃宁县| 华蓥市| 丹凤县| 榆社县| 光泽县| 开封市| 榆中县| 石景山区| 泾阳县| 读书| 巴楚县| 开鲁县| 鄂托克旗| 福清市| 东乌珠穆沁旗| 诏安县| 平陆县| 会宁县| 江油市|