Spring Boot中的WebSocket實時通信實戰(zhàn)解析
WebSocket 在 Spring Boot 中的實戰(zhàn)解析:實時通信的技術(shù)利器
一、引言:為什么我們需要 WebSocket?
在傳統(tǒng)的 Web 應(yīng)用中,客戶端(瀏覽器)與服務(wù)器之間的通信是 請求-響應(yīng) 模式:客戶端發(fā)起請求,服務(wù)器處理后返回結(jié)果。這種模式適用于大多數(shù)場景,但在需要 實時雙向通信 的場景下(如聊天室、股票行情、在線協(xié)作、游戲等),頻繁輪詢(Polling)或長輪詢(Long Polling)會帶來高延遲、高開銷的問題。
WebSocket 協(xié)議應(yīng)運而生——它提供了一種全雙工、低延遲、持久化的通信通道,允許服務(wù)器主動向客戶端推送數(shù)據(jù),徹底改變了 Web 實時交互的格局。
在 Spring Boot 生態(tài)中,通過 spring-boot-starter-websocket 模塊,我們可以輕松集成 WebSocket,并結(jié)合 STOMP(Simple Text-Oriented Messaging Protocol)實現(xiàn)更高級的消息路由與訂閱機制。
1.1 從“輪詢”到“長連接”的進化史
在WebSocket出現(xiàn)之前,實現(xiàn)實時通信主要靠這些“土辦法”:
| 技術(shù)方案 | 工作原理 | 缺點 |
|---|---|---|
| 短輪詢 | 客戶端每隔幾秒問一次:“有新消息嗎?” | 浪費帶寬,實時性差 |
| 長輪詢 | 客戶端問“有新消息嗎?”,服務(wù)器hold住,有消息才回復(fù) | 連接占用時間長,服務(wù)器壓力大 |
| SSE | 服務(wù)器單向推送,客戶端只能接收 | 單向通信,功能有限 |
WebSocket的登場改變了游戲規(guī)則:
- 一次握手,持久連接:建立連接后,雙向通道一直打開
- 服務(wù)端主動推送:服務(wù)器想什么時候發(fā)就什么時候發(fā)
- 極低的通信開銷:沒有HTTP頭部的重復(fù)傳輸
1.2 WebSocket vs HTTP:本質(zhì)區(qū)別
HTTP交互流程(像發(fā)短信):

每次請求都要重復(fù):建立TCP連接 → TLS握手 → 發(fā)送HTTP頭部 → 傳輸數(shù)據(jù) → 斷開連接
WebSocket交互流程(像打電話):

二者的關(guān)鍵區(qū)別:

二、Spring Boot 中的 WebSocket 技術(shù)棧
Spring 對 WebSocket 的支持分為兩個層次:
| 層級 | 技術(shù) | 說明 |
|---|---|---|
| 底層 | javax.websocket 或 Spring 原生 WebSocket | 直接處理原始 WebSocket 消息 |
| 高層 | STOMP over WebSocket | 基于消息代理的發(fā)布/訂閱模型,更易開發(fā) |
? 推薦使用 STOMP:它提供了類似 JMS 的語義(目的地、訂閱、廣播),適合復(fù)雜業(yè)務(wù)場景。
2.1 核心注解與類
| 組件 | 作用 |
|---|---|
@EnableWebSocketMessageBroker | 啟用 WebSocket 消息代理 |
WebSocketMessageBrokerConfigurer | 配置 STOMP 端點與消息代理 |
@MessageMapping | 映射客戶端發(fā)送到特定路徑的消息 |
SimpMessagingTemplate | 服務(wù)端主動向客戶端推送消息 |
2.2 三步快速集成
第一步:添加依賴
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>第二步:配置WebSocket
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
// 注冊處理器,指定連接路徑
registry.addHandler(myWebSocketHandler(), "/ws")
.setAllowedOrigins("*") // 生產(chǎn)環(huán)境要限制具體域名
.withSockJS(); // 為不支持WebSocket的瀏覽器提供降級方案
}
@Bean
public WebSocketHandler myWebSocketHandler() {
return new MyWebSocketHandler();
}
}第三步:實現(xiàn)核心處理器
@Component
public class MyWebSocketHandler extends TextWebSocketHandler {
// 保存所有活躍連接
private static final Map<String, WebSocketSession> sessions =
new ConcurrentHashMap<>();
@Override
public void afterConnectionEstablished(WebSocketSession session) {
// 連接建立時調(diào)用
String userId = extractUserId(session);
sessions.put(userId, session);
log.info("用戶 {} 連接成功,當(dāng)前在線: {} 人", userId, sessions.size());
// 發(fā)送歡迎消息
session.sendMessage(new TextMessage("連接成功!"));
}
@Override
protected void handleTextMessage(WebSocketSession session,
TextMessage message) {
// 處理客戶端發(fā)送的消息
String payload = message.getPayload();
log.info("收到消息: {}", payload);
// 處理業(yè)務(wù)邏輯,比如廣播或定向回復(fù)
handleMessage(session, payload);
}
@Override
public void afterConnectionClosed(WebSocketSession session,
CloseStatus status) {
// 連接關(guān)閉時調(diào)用
String userId = extractUserId(session);
sessions.remove(userId);
log.info("用戶 {} 斷開連接,原因: {}", userId, status);
}
}2.3 進階:使用STOMP協(xié)議
2.3.1 什么是STOMP?
STOMP(Simple Text Oriented Messaging Protocol)是一個簡單的文本協(xié)議,它為WebSocket提供了更高級的消息模式。如果說原始的WebSocket是"裸奔",那么STOMP就是給它穿上了"協(xié)議的外衣"。
STOMP的核心概念:
- Destination(目的地):消息發(fā)送的目標(biāo)地址
- SUBSCRIBE(訂閱):客戶端訂閱某個目的地
- SEND(發(fā)送):客戶端向目的地發(fā)送消息
- MESSAGE(消息):服務(wù)器向客戶端推送消息
2.3.2 STOMP協(xié)議結(jié)構(gòu)
一個簡單的STOMP幀示例:
SEND
destination:/app/chat
content-type:application/json
content-length:23
{"text":"Hello World!"}響應(yīng)幀:
MESSAGE
destination:/topic/chat
content-type:application/json
content-length:45
subscription:sub-0
message-id:msg-123
{"user":"Tom","text":"Hello World!"}2.3.3 SpringBoot中的STOMP配置
@Configuration
@EnableWebSocketMessageBroker // 啟用STOMP消息代理
public class WebSocketStompConfig implements WebSocketMessageBrokerConfigurer {
@Override
public void registerStompEndpoints(StompEndpointRegistry registry) {
// 注冊STOMP端點
registry.addEndpoint("/ws-stomp")
.setAllowedOrigins("*")
.withSockJS(); // SockJS支持
}
@Override
public void configureMessageBroker(MessageBrokerRegistry registry) {
// 配置消息代理
registry.enableSimpleBroker("/topic", "/queue"); // 客戶端訂閱的前綴
registry.setApplicationDestinationPrefixes("/app"); // 客戶端發(fā)送消息的前綴
registry.setUserDestinationPrefix("/user"); // 用戶私信前綴
}
}2.3.4 STOMP控制器示例
@Controller
public class ChatController {
// 處理發(fā)送到/app/chat的消息
@MessageMapping("/chat")
@SendTo("/topic/chat") // 廣播給所有訂閱/topic/chat的客戶端
public ChatMessage handleChatMessage(ChatMessage message) {
message.setTimestamp(LocalDateTime.now());
log.info("收到聊天消息: {}", message);
return message;
}
// 處理私信
@MessageMapping("/private")
public void sendPrivateMessage(@Payload PrivateMessage message,
SimpMessageHeaderAccessor headerAccessor) {
// 獲取發(fā)送者
String sender = headerAccessor.getUser().getName();
// 使用convertAndSendToUser發(fā)送給特定用戶
simpMessagingTemplate.convertAndSendToUser(
message.getRecipient(), // 接收者用戶名
"/queue/private", // 用戶私信隊列
new PrivateMessage(sender, message.getRecipient(), message.getContent())
);
}
// 處理訂閱通知
@EventListener
public void handleSessionSubscribe(SessionSubscribeEvent event) {
StompHeaderAccessor headers = StompHeaderAccessor.wrap(event.getMessage());
String destination = headers.getDestination();
String sessionId = headers.getSessionId();
log.info("Session {} 訂閱了 {}", sessionId, destination);
}
}2.4.5 前端STOMP客戶端示例
// 連接STOMP服務(wù)器
const socket = new SockJS('/ws-stomp');
const stompClient = Stomp.over(socket);
// 連接成功回調(diào)
stompClient.connect({}, function(frame) {
console.log('Connected: ' + frame);
// 訂閱公共聊天頻道
stompClient.subscribe('/topic/chat', function(message) {
const chatMsg = JSON.parse(message.body);
showMessage(chatMsg);
});
// 訂閱個人私信隊列
stompClient.subscribe('/user/queue/private', function(message) {
const privateMsg = JSON.parse(message.body);
showPrivateMessage(privateMsg);
});
});
// 發(fā)送聊天消息
function sendMessage() {
const message = {
content: document.getElementById('message').value,
sender: currentUser
};
stompClient.send("/app/chat", {}, JSON.stringify(message));
}
// 發(fā)送私信
function sendPrivateMessage(toUser, content) {
const message = {
recipient: toUser,
content: content
};
stompClient.send("/app/private", {}, JSON.stringify(message));
}三、WebSocket的核心技術(shù)要點
3.1 連接生命周期管理
連接生命周期包括四個關(guān)鍵階段,每個階段需要處理不同的業(yè)務(wù)邏輯:
建立連接時(afterConnectionEstablished)需要處理:
- 驗證用戶身份 - 檢查token或session,確保連接合法性
- 初始化會話狀態(tài) - 創(chuàng)建用戶會話上下文,保存必要信息
- 通知相關(guān)服務(wù)用戶上線 - 更新用戶在線狀態(tài),通知好友
- 發(fā)送未讀消息 - 推送離線期間積累的消息
@Component
public class WebSocketLifecycleManager {
public void onOpen(String sessionId, String userId) {
// 1. 驗證用戶身份
if (!userService.validateToken(userId, getToken(sessionId))) {
closeConnection(sessionId, CloseStatus.NOT_ACCEPTABLE);
return;
}
// 2. 初始化會話狀態(tài)
SessionContext context = new SessionContext(userId, sessionId);
sessionStore.save(sessionId, context);
// 3. 通知相關(guān)服務(wù)用戶上線
presenceService.userOnline(userId);
notifyFriends(userId, true); // 通知好友用戶上線
// 4. 發(fā)送未讀消息
List<Message> unreadMessages = messageService.getUnreadMessages(userId);
unreadMessages.forEach(msg -> sendMessage(sessionId, msg));
}
}消息處理時(handleTextMessage)需要處理:
- 消息格式驗證 - 檢查JSON格式、必要字段
- 業(yè)務(wù)邏輯處理 - 根據(jù)消息類型執(zhí)行不同業(yè)務(wù)
- 消息持久化 - 保存到數(shù)據(jù)庫,確保不丟失
- 響應(yīng)或轉(zhuǎn)發(fā) - 回復(fù)發(fā)送者或轉(zhuǎn)發(fā)給其他用戶
public void onMessage(String sessionId, String rawMessage) {
// 1. 消息格式驗證
Message message;
try {
message = jsonMapper.readValue(rawMessage, Message.class);
validateMessage(message);
} catch (Exception e) {
sendError(sessionId, "消息格式錯誤");
return;
}
// 2. 業(yè)務(wù)邏輯處理
switch (message.getType()) {
case "CHAT":
handleChatMessage(sessionId, message);
break;
case "COMMAND":
handleCommand(sessionId, message);
break;
// ... 其他消息類型
}
// 3. 消息持久化
messageService.saveMessage(message);
// 4. 響應(yīng)或轉(zhuǎn)發(fā)
if (message.needResponse()) {
sendResponse(sessionId, createResponse(message));
}
if (message.needForward()) {
forwardMessage(message.getTarget(), message);
}
}連接關(guān)閉時(afterConnectionClosed)需要處理:
- 清理會話資源 - 釋放內(nèi)存,關(guān)閉相關(guān)資源
- 更新用戶狀態(tài)為離線 - 標(biāo)記用戶下線時間
- 記錄斷開原因 - 用于分析連接穩(wěn)定性
- 通知相關(guān)服務(wù) - 通知好友用戶下線
public void onClose(String sessionId, CloseStatus status) {
// 1. 清理會話資源
SessionContext context = sessionStore.remove(sessionId);
if (context != null) {
context.cleanup();
}
// 2. 更新用戶狀態(tài)為離線
String userId = getUserIdFromSession(sessionId);
presenceService.userOffline(userId);
// 3. 記錄斷開原因
connectionLogService.logDisconnect(sessionId, userId, status.getCode(), status.getReason());
// 4. 通知相關(guān)服務(wù)
notifyFriends(userId, false); // 通知好友用戶下線
cleanupUserSubscriptions(userId); // 清理用戶的所有訂閱
}錯誤時(handleTransportError)需要處理:
- 記錄錯誤日志 - 詳細記錄異常信息
- 嘗試恢復(fù)連接 - 對于可恢復(fù)錯誤嘗試重連
- 通知監(jiān)控系統(tǒng) - 觸發(fā)告警,人工干預(yù)
- 優(yōu)雅降級 - 切換到備用通信方式
public void onError(String sessionId, Throwable error) {
// 1. 記錄錯誤日志
log.error("WebSocket連接錯誤 sessionId: {}", sessionId, error);
// 2. 嘗試恢復(fù)連接(如果是網(wǎng)絡(luò)波動等臨時錯誤)
if (isRecoverableError(error)) {
scheduleReconnection(sessionId);
} else {
// 3. 通知監(jiān)控系統(tǒng)
alertService.sendAlert("WebSocket連接異常",
"Session: " + sessionId + ", Error: " + error.getMessage());
// 4. 優(yōu)雅降級
fallbackToHttp(sessionId, getUserIdFromSession(sessionId));
closeConnection(sessionId, CloseStatus.SERVER_ERROR);
}
}3.2 心跳機制與健康檢查
@Configuration
public class HeartbeatConfig {
@Bean
public ServletServerContainerFactoryBean createWebSocketContainer() {
ServletServerContainerFactoryBean container =
new ServletServerContainerFactoryBean();
// 重要配置
container.setMaxSessionIdleTimeout(300000L); // 5分鐘無活動斷開
container.setMaxTextMessageBufferSize(8192); // 最大消息大小
container.setMaxBinaryMessageBufferSize(8192);
container.setAsyncSendTimeout(5000L); // 異步發(fā)送超時
return container;
}
}
// 心跳檢測實現(xiàn)
@Component
public class HeartbeatService {
private final ScheduledExecutorService scheduler =
Executors.newScheduledThreadPool(1);
@PostConstruct
public void startHeartbeat() {
scheduler.scheduleAtFixedRate(() -> {
checkConnections();
sendPing();
}, 30, 30, TimeUnit.SECONDS); // 每30秒檢測一次
}
private void checkConnections() {
// 檢查所有連接的健康狀態(tài)
// 移除僵尸連接
// 記錄連接統(tǒng)計信息
}
private void sendPing() {
// 向所有活躍連接發(fā)送ping消息
// 處理未響應(yīng)pong的連接
}
}3.3 消息可靠性與重連機制
public class ReliableMessageService {
// 消息確認機制
public void sendWithAck(String sessionId, String message) {
String msgId = generateMsgId();
// 發(fā)送消息
webSocketHandler.send(sessionId, wrapMessage(msgId, message));
// 啟動確認計時器
scheduler.schedule(() -> {
if (!isAcked(msgId)) {
log.warn("消息 {} 未確認,嘗試重發(fā)", msgId);
retrySend(sessionId, msgId, message);
}
}, 5, TimeUnit.SECONDS);
}
// 客戶端重連處理
public void handleReconnect(String oldSessionId, String newSessionId) {
// 1. 轉(zhuǎn)移會話狀態(tài)
// 2. 重發(fā)未確認消息
// 3. 恢復(fù)訂閱關(guān)系
// 4. 更新會話映射
}
// 消息去重
public boolean isDuplicate(String msgId) {
// 基于Redis或本地緩存實現(xiàn)
// 防止重復(fù)處理消息
return false;
}
}四、實戰(zhàn):寫一個監(jiān)控SpringBoot應(yīng)用的實時監(jiān)控告警系統(tǒng)
Spring Insight 是我的一個開源項目,目前正在緊張的開發(fā)中。項目地址:https://github.com/iweidujiang/spring-insight,歡迎關(guān)注,順便求個 star,哈哈。
在監(jiān)控診斷類工具中,WebSocket 可以:
- 實時告警:第一時間發(fā)現(xiàn)問題
- 動態(tài)拓撲:實時展示微服務(wù)依賴變化
- 性能監(jiān)控:實時推送指標(biāo)數(shù)據(jù)
- 在線診斷:實時查看日志和跟蹤信息
例,實時統(tǒng)計當(dāng)前監(jiān)控信息:
/**
* 廣播實時統(tǒng)計信息(每5秒一次)
*/
@Scheduled(fixedDelay = 5000)
public void broadcastStats() {
if (connectionCount.get() == 0) return;
try {
// 獲取最新數(shù)據(jù)
var collectorStats = dataCollectorService.getCollectorStats();
var serviceStats = dataCollectorService.getServiceStats();
var errorAnalysis = dataCollectorService.getErrorAnalysis(1); // 最近1小時
// 構(gòu)建消息
Map<String, Object> data = new HashMap<>();
data.put("collectorStats", collectorStats);
data.put("serviceStats", serviceStats.subList(0, Math.min(5, serviceStats.size())));
data.put("errorAnalysis", errorAnalysis.subList(0, Math.min(5, errorAnalysis.size())));
data.put("timestamp", Instant.now().toString());
data.put("cacheSize", dataCollectorService.getCacheSize());
WebSocketMessage message = new WebSocketMessage();
message.setType("STATS_UPDATE");
message.setData(data);
// 廣播消息
messagingTemplate.convertAndSend("/topic/stats", message);
log.debug("廣播實時統(tǒng)計信息");
} catch (Exception e) {
log.error("廣播實時統(tǒng)計信息失敗", e);
}
}五、其他典型使用場景
| 場景 | 說明 | WebSocket 優(yōu)勢 |
|---|---|---|
| 在線客服/聊天系統(tǒng) | 用戶與客服實時對話 | 低延遲、支持多房間 |
| 股票/金融行情推送 | 實時價格更新 | 減少服務(wù)器壓力,避免輪詢 |
| 協(xié)同編輯 | 多人同時編輯文檔 | 實時同步操作,沖突檢測 |
| 游戲狀態(tài)同步 | 多人在線小游戲 | 高頻消息傳遞,毫秒級響應(yīng) |
| IoT 設(shè)備監(jiān)控 | 傳感器數(shù)據(jù)上報 | 長連接節(jié)省資源 |
? 用WebSocket:
- 實時雙向通信需求(聊天、協(xié)作)
- 高頻數(shù)據(jù)推送(監(jiān)控、行情)
- 低延遲要求(游戲、實時控制)
- 服務(wù)端主動通知(告警、狀態(tài)更新)
? 不用WebSocket:
- 簡單的請求-響應(yīng)模式(REST API足夠)
- 客戶端偶爾拉取數(shù)據(jù)(用HTTP輪詢)
- 單向信息流(考慮SSE)
- 移動端弱網(wǎng)絡(luò)環(huán)境(可能連接不穩(wěn)定)
六、總結(jié)
WebSocket 是構(gòu)建現(xiàn)代實時 Web 應(yīng)用的基石。Spring Boot 通過簡潔的配置和強大的 STOMP 支持,讓開發(fā)者能夠快速實現(xiàn)高性能、可擴展的雙向通信系統(tǒng)。
記住:不是所有場景都需要 WebSocket。對于低頻更新(如每分鐘一次),傳統(tǒng) REST + 定時輪詢可能更簡單。但在高頻、低延遲、事件驅(qū)動的場景下,WebSocket 幾乎是唯一選擇。
到此這篇關(guān)于WebSocket 在 Spring Boot 中的實戰(zhàn)解析:實時通信的技術(shù)利器的文章就介紹到這了,更多相關(guān)Spring Boot WebSocket 實時通信內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
從基礎(chǔ)到高級詳解Java讀寫Excel公式的實戰(zhàn)指南
這篇文章主要為大家詳細介紹了Java使用Spireire.XLSforJava處理Excel公式的方法,涵蓋寫入、讀取、跨工作表引用、日期時間函數(shù)等操作,希望對大家有所幫助2026-05-05
Spring解決循環(huán)依賴的方法及三級緩存機制實踐案例
Spring通過三級緩存解決單例Bean循環(huán)依賴,但無法處理構(gòu)造器、prototype作用域及@Async場景,建議使用setter注入、@Lazy注解和架構(gòu)優(yōu)化,遵循設(shè)計原則避免依賴問題,本文介紹Spring如何解決循環(huán)依賴:深入理解三級緩存機制,感興趣的朋友一起看看吧2025-09-09
java使用異或?qū)崿F(xiàn)變量互換和異或加密解密示例
這篇文章主要介紹了使用異或?qū)崿F(xiàn)變量互換和異或加密解密示例,需要的朋友可以參考下2014-02-02
詳解FileInputStream讀取文件數(shù)據(jù)的兩種方式
這篇文章主要介紹了詳解FileInputStream讀取文件數(shù)據(jù)的兩種方式,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-08-08

