WebSocket(java版)服務(wù)核心實例代碼
說明:
這是一個使用 Java JDK 8 和 Spring Boot 實現(xiàn)的WebSocket演示項目。目的是為解決多端消息通訊的問題。
WebSocket 是一種基于 TCP 的全雙工通信協(xié)議,核心作用是解決傳統(tǒng) HTTP 協(xié)議 “請求 - 響應(yīng)” 模式的局限性,實現(xiàn) 客戶端與服務(wù)器之間的實時、雙向、低延遲數(shù)據(jù)傳輸。
源碼地址:https://gitee.com/lqh4188/web-socket
一、功能介紹
功能特性:
- 基于 Maven 的 Spring Boot 項目骨架。
- 純 WebSocket 端點 /ws ,支持用戶隔離,http:使用ws,https:使用wss。
- 支持分片設(shè)置和緩沖區(qū)大小設(shè)置,解決傳輸內(nèi)容限制
- 提供靜態(tài)測試頁面 index.html ,用于連接、發(fā)送消息、查看消息。
項目結(jié)構(gòu):
- pom.xml :Spring Boot 3.3,依賴 spring-boot-starter-web 和 spring-boot-starter-websocket 。
- src/main/java/com/example/websocket/WebSocketApplication.java :應(yīng)用入口。
- src/main/java/com/example/websocket/WebSocketConfig.java :注冊 WebSocket 處理器,端點為 /ws 。
- src/main/java/com/example/websocket/ChatWebSocketHandler.java :文本消息處理,廣播到所有會話。
- src/main/resources/static/index.html :頁面內(nèi)置 JS,連接 ws://{host}/ws ,可發(fā)送、顯示消息。
關(guān)鍵代碼位置
- 啟動類: src/main/java/com/example/websocket/WebSocketApplication.java:1
- WebSocket 配置: src/main/java/com/example/websocket/WebSocketConfig.java:1
- 文本消息處理器: src/main/java/com/example/websocket/ChatWebSocketHandler.java:1
- 靜態(tài)頁面: src/main/resources/static/index.html:1
測試連接
- 打開 http://localhost:8800 ,使用頁面上的“連接/發(fā)送”測試
- WebSocket 地址: ws://localhost:8080/ws
二、運行測試
可通過UserId來創(chuàng)建獨立的聯(lián)接,進行用戶隔離

三、核心代碼說明
由于websocket對傳輸?shù)膬?nèi)容有限制,若內(nèi)容較大可進行緩沖區(qū)大小設(shè)置,并對不同文本進行分片處理
ChatWebSocketHandler.java代碼:
package com.example.websocket;
import java.io.ByteArrayOutputStream;
import java.net.URI;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import org.springframework.web.socket.BinaryMessage;
import org.springframework.web.socket.CloseStatus;
import org.springframework.web.socket.TextMessage;
import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.handler.AbstractWebSocketHandler;
import com.fasterxml.jackson.databind.ObjectMapper;
public class ChatWebSocketHandler extends AbstractWebSocketHandler {
private final ConcurrentHashMap<String, Set<WebSocketSession>> userSessions = new ConcurrentHashMap<>();
private static final ObjectMapper MAPPER = new ObjectMapper();
private final ConcurrentHashMap<String, StringBuilder> textFragments = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, ByteArrayOutputStream> binaryFragments = new ConcurrentHashMap<>();
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
// 驗證用戶ID的有效性
String uid = resolveUserId(session);
if (uid == null || uid.isEmpty()) {
session.close(CloseStatus.BAD_DATA);
return;
}
session.getAttributes().put("userId", uid);
//多會話管理
userSessions.computeIfAbsent(uid, k -> ConcurrentHashMap.newKeySet()).add(session);
}
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
// 分片處理
String id = session.getId();
if (!message.isLast()) {
textFragments.computeIfAbsent(id, k -> new StringBuilder()).append(message.getPayload());
return;
}
StringBuilder sb = textFragments.remove(id);
String payload = sb != null ? sb.append(message.getPayload()).toString() : message.getPayload();
routePayload(session, payload);
}
@Override
protected void handleBinaryMessage(WebSocketSession session, BinaryMessage message) throws Exception {
//二進制消息處理
String id = session.getId();
ByteBuffer buf = message.getPayload();
byte[] chunk = new byte[buf.remaining()];
buf.get(chunk);
ByteArrayOutputStream acc = binaryFragments.computeIfAbsent(id, k -> new ByteArrayOutputStream());
acc.write(chunk);
if (message.isLast()) {
byte[] all = acc.toByteArray();
binaryFragments.remove(id);
String payload = new String(all, StandardCharsets.UTF_8);
routePayload(session, payload);
}
}
@Override
public boolean supportsPartialMessages() {
return true;
}
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception {
// WebSocket 連接關(guān)閉時的清理邏輯
Object v = session.getAttributes().get("userId");
if (v == null) return;
String uid = String.valueOf(v);
Set<WebSocketSession> set = userSessions.get(uid);
if (set != null) {
set.remove(session);
if (set.isEmpty()) userSessions.remove(uid);
}
}
/** 從 WebSocket 連接的 URL 查詢參數(shù)中提取用戶ID */
private String resolveUserId(WebSocketSession session) {
URI uri = session.getUri();
if (uri == null) return null;
String q = uri.getQuery();
if (q == null || q.isEmpty()) return null;
String[] parts = q.split("&");
for (String p : parts) {
int i = p.indexOf('=');
if (i > 0) {
String k = p.substring(0, i);
String val = p.substring(i + 1);
if ("userId".equals(k)) return val;
}
}
return null;
}
private void routePayload(WebSocketSession session, String payload) throws Exception {
Object v = session.getAttributes().get("userId");
if (v == null) return;
String fromUid = String.valueOf(v);
// 解析消息
Message message = new Message();
message.setFromUserId(fromUid);
try {
// 嘗試將payload解析為Message對象
Message receivedMsg = MAPPER.readValue(payload, Message.class);
message.setToUserId(receivedMsg.getToUserId());
message.setContent(receivedMsg.getContent());
message.setType(receivedMsg.getType());
} catch (Exception e) {
// 如果解析失敗,將整個payload作為content
message.setContent(payload);
}
String toUid = message.getToUserId();
boolean isP2P = toUid != null && !toUid.isEmpty();
Set<WebSocketSession> targets;
if (isP2P) {
targets = userSessions.get(toUid);
} else {
targets = userSessions.get(fromUid);
}
// 序列化消息對象
String outStr = MAPPER.writeValueAsString(message);
TextMessage msg = new TextMessage(outStr);
if (targets == null || targets.isEmpty()) {
if (session.isOpen()) {
session.sendMessage(msg);
}
return;
}
for (WebSocketSession s : targets) {
if (s.isOpen()) {
s.sendMessage(msg);
}
}
if (isP2P && session.isOpen()) {
session.sendMessage(msg);
}
}
}
配置類WebSocketConfig.java
package com.example.websocket;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.WebSocketHandler;
import org.springframework.web.socket.config.annotation.EnableWebSocket;
import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry;
import org.springframework.web.socket.server.standard.ServletServerContainerFactoryBean;
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(chatHandler(), "/ws").setAllowedOriginPatterns("*");
}
@Bean
public WebSocketHandler chatHandler() {
return new ChatWebSocketHandler();
}
// 配置 WebSocket 容器參數(shù)(解決消息過大、超時等問題)
@Bean
public ServletServerContainerFactoryBean createWebSocketContainer() {
ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean();
// 文本消息緩沖區(qū):2MB(解決解碼后消息過大的核心配置)
container.setMaxTextMessageBufferSize(2 * 1024 * 1024);
// 二進制消息緩沖區(qū):4MB(按需配置)
container.setMaxBinaryMessageBufferSize(4 * 1024 * 1024);
// 會話空閑超時:60秒(無交互則關(guān)閉連接)
container.setMaxSessionIdleTimeout(60_000L);
return container;
}
}
總結(jié)
到此這篇關(guān)于WebSocket(java版)服務(wù)核心代碼的文章就介紹到這了,更多相關(guān)java WebSocket服務(wù)內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- Java Socket編程實現(xiàn)群聊實踐案例
- Java通過ServerSocket與Socket實現(xiàn)通信過程
- java網(wǎng)絡(luò)編程之socket網(wǎng)絡(luò)編程示例(服務(wù)器端/客戶端)
- 使用Java和WebSocket實現(xiàn)網(wǎng)頁聊天室實例代碼
- java使用Socket類接收和發(fā)送數(shù)據(jù)
- java socket編程實例代碼講解
- 簡單的java socket客戶端和服務(wù)端示例
- Java后端Tomcat實現(xiàn)WebSocket實例教程
- 教你怎么使用Java實現(xiàn)WebSocket
- Java Socket通信(一)之客戶端程序 發(fā)送和接收數(shù)據(jù)
- 基于Java Socket實現(xiàn)一個簡易在線聊天功能(一)
- Java網(wǎng)絡(luò)編程:Socket套接字、TCP和UDP協(xié)議以及Java高性能實戰(zhàn)
相關(guān)文章
詳解Springboot2.3集成Spring security 框架(原生集成)
這篇文章主要介紹了詳解Springboot2.3集成Spring security 框架(原生集成),文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-08-08
Elasticsearch?Recovery索引分片分配詳解
這篇文章主要為大家介紹了關(guān)于Elasticsearch的Recovery索引分片分配詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪<BR>2022-04-04

