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

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

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

引言

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

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

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

一、引入依賴

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是整個(gè)方案的核心。這里有幾個(gè)設(shè)計(jì)點(diǎn)值得注意:

  • 支持多標(biāo)簽頁(yè):一個(gè)用戶可能打開多個(gè)瀏覽器標(biāo)簽頁(yè),所以用Map<String, Map<String, SseEmitter>>來組織,外層key是userId,內(nèi)層key是token(每個(gè)標(biāo)簽頁(yè)唯一)。
  • 生命周期管理:通過onCompletion、onTimeout、onError回調(diào)來清理連接,防止內(nèi)存泄漏。
  • 不超時(shí)設(shè)置:new SseEmitter(0L)表示永不超時(shí),你也可以根據(jù)業(yè)務(wù)需要設(shè)置一個(gè)合理的超時(shí)時(shí)間(如30分鐘)。
  • 全站廣播:broadcast()方法會(huì)遍歷所有在線用戶并推送,適合系統(tǒng)公告類消息。
@Component
@Slf4j
public class SseEmitterManager {
    /**
     * 用戶ID -> (連接token -> SseEmitter)
     * 一個(gè)用戶可能有多個(gè)瀏覽器標(biāo)簽頁(yè),用token區(qū)分
     */
    private final Map<String, Map<String, SseEmitter>> userEmitters = new ConcurrentHashMap<>();
    /**
     * 建立SSE連接
     * @param userId 用戶ID
     * @param token 連接標(biāo)識(shí)(可用UUID)
     */
    public SseEmitter connect(String userId, String token) {
        // 超時(shí)時(shí)間設(shè)為0表示不超時(shí),也可設(shè)置具體毫秒數(shù)
        SseEmitter emitter = new SseEmitter(0L);
        // 注冊(cè)回調(diào):連接關(guān)閉時(shí)清理
        emitter.onCompletion(() -> removeEmitter(userId, token));
        emitter.onTimeout(() -> removeEmitter(userId, token));
        emitter.onError(e -> removeEmitter(userId, token));
        // 存儲(chǔ)連接
        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());
    }
    /**
     * 獲取當(dāng)前在線人數(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實(shí)現(xiàn)了MessageListener接口,它會(huì)監(jiān)聽station:message頻道。當(dāng)收到Redis消息時(shí):

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

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

@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提供的消息容器,它會(huì)自動(dòng)管理訂閱和監(jiān)聽線程。注意這里訂閱的是station:message頻道,你可以根據(jù)業(yè)務(wù)需要定義多個(gè)頻道(比如station:notice、station:system)。

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

@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;
    }
}

五、消息實(shí)體

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

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

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

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

這種設(shè)計(jì)的優(yōu)點(diǎn)是:發(fā)送方不需要知道消息最終由哪個(gè)實(shí)例處理,也不需要維護(hù)實(shí)例之間的網(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,所有實(shí)例都會(huì)收到
            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對(duì)外暴露了三個(gè)核心接口:

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

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

@RestController
@RequestMapping("/api/sse")
@Slf4j
public class SseController {
    
    @Autowired
    private SseEmitterManager sseEmitterManager;
    
    @Autowired
    private MessageService messageService;
    
    /**
     * SSE連接端點(diǎn)
     */
    @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實(shí)現(xiàn)跨實(shí)例SSE消息推送的詳細(xì)內(nèi)容,更多關(guān)于SpringBoot跨實(shí)例SSE消息推送的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

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

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

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

    淺談JVM中的JOL

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

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

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

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

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

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

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

    JAVA--HashMap熱門面試題

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

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

    這篇文章主要介紹了JDK動(dòng)態(tài)代理接口和接口實(shí)現(xiàn)類,JDK動(dòng)態(tài)代理是代理模式的一種實(shí)現(xiàn)方式,因?yàn)樗腔诮涌趤碜龃淼?所以也常被稱為接口代理,文中通過實(shí)例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-06-06
  • 淺談@PostConstruct不被調(diào)用的原因

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

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

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

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

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

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

最新評(píng)論

九龙县| 札达县| 都江堰市| 红桥区| 永春县| 威远县| 怀安县| 满洲里市| 建宁县| 怀安县| 右玉县| 东山县| 高淳县| 石城县| 武强县| 尖扎县| 铜鼓县| 元江| 鲁甸县| 荣成市| 唐山市| 华容县| 新沂市| 新龙县| 昭苏县| 新龙县| 墨江| 星子县| 伊宁市| 翁牛特旗| 聊城市| 信宜市| 克拉玛依市| 镇平县| 清徐县| 邵阳县| 湖南省| 禄劝| 若羌县| 台东市| 辰溪县|