websocket消息传输
全双工通信:WebSocket 协议建立在 TCP 连接上,支持客户端和服务器之间的实时双向消息传输,适合实时通讯和即时更新的应用。
一、websocket实现前后端通讯
1.引入依赖,springboot支持websocket
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>
2.注意使用websocket首先要进行bean的配置开启对websocket的支持
@Configuration
public class WebSocketConfig {
/**
* 开启webSocket服务
* @return
*/
@Bean
public ServerEndpointExporter serverEndpointExporter(){
return new ServerEndpointExporter();
}
}
3.配置服务类
实现websocket通讯的类
每建立一个新的连接,都会创建一个websocket的对象,因此每次连接后的对象的数据都是独属于本次webcocket连接到的对象,除非设置为静态变量作为公用的变量。
//设置websocket连接时的url
@ServerEndpoint("/websocket/{sid}")
//将此对象的创建交给springboot
@Component
public class WebSocketServer {
/*变量*/
private String sid = "";
private Session session;
......
/*连接成功调用的方法*/
通过OnOpen注解来标注
@OnOpen
public void onOpen(Session session, @PathParam("sid") String param) {
......
}
/*连接关闭调用方法*/
通过OnClose注解来标注
@OnClose
public void onClose() {
......
}
/*接收客户端发送的信息,并进行转发*/
通过OnMessage注解来标注
@OnMessage
public void onMessage(String message, Session session) throws IOException {
......
}
/*出现异常时调用方法*/
通过OnError注解来标注
@OnError
public void onError(Session session, Throwable throwable) {
......
}
/*其他方法*/
}
3.1变量
1.CopyOnWriteArraySet
CopyOnWriteArraySet 是 Java 中的一个线程安全集合类,它是 Set 接口的一种实现。它的特点是使用“写时复制”(Copy-On-Write)的机制来保证线程安全性。
使用此集合来存储websocket对象。
此变量可以用于群发。
private static CopyOnWriteArraySet<WebSocketServer> webSocketSet = new CopyOnWriteArraySet<>();
群发方法
/**
* 群发自定义消息
*/
public static void sendInfo(String message, @PathParam("sid") String sid) throws IOException {
log.info(" 将消息推送给 " + sid + " ,推送内容为 :" + message);
for (WebSocketServer item : webSocketSet) {
try {
// 这里可以设定只推送给这个 sid 的,为 null 则全部推送
if (sid == null) {
item.sendMessage(message);
} else if (item.sid.equals(sid)) {
item.sendMessage(message);
}
} catch (IOException e) {
continue;
}
}
}
2.AtomicInteger
private static AtomicInteger onlineCount = new AtomicInteger();
AtomicInteger 是 Java 中
java.util.concurrent.atomic 包下的一个类,用于在高并发环境下实现对整数值的原子操作。它提供了一些方法来确保对整数值的操作是以原子方式执行的,从而避免了在多线程环境下可能出现的竞争条件和数据不一致问题。
AtomicInteger 的主要特点包括:
- 原子性操作
AtomicInteger 通过使用硬件级别的原子指令(如 CAS,Compare-And-Swap)来确保对整数值的操作是原子的,避免了传统同步机制(如 synchronized 和 ReentrantLock)带来的性能开销。
- 无锁实现
由于使用了 CAS 操作,AtomicInteger 在多数情况下无需使用锁,即使在高并发场景下也能提供良好的性能。
- 线程安全
所有对 AtomicInteger 的操作都是线程安全的,适合在多线程环境中使用。
AtomicInteger 提供了一系列原子操作的方法,例如:
- get(): 获取当前值。
- set(int newValue): 设置为指定的新值。
- getAndSet(int newValue): 获取当前值,并设置为指定的新值。
- compareAndSet(int expect, int update): 如果当前值等于预期值,则设置为新值。
- incrementAndGet(): 原子地将当前值加 1,并返回更新后的值。
- getAndIncrement(): 原子地将当前值加 1,并返回更新前的值。
- decrementAndGet(): 原子地将当前值减 1,并返回更新后的值。
- getAndDecrement(): 原子地将当前值减 1,并返回更新前的值。
- addAndGet(int delta): 原子地将当前值加上指定的增量,并返回更新后的值。
- getAndAdd(int delta): 原子地将当前值加上指定的增量,并返回更新前的值。
3.Session
private Session session;
此Session是websocket中的重要部分,里面存放了客户端连接时的一些信息

获取客户端ip可以见下方博客,适用于jdk8版本,jdk17继承类改变该方法不适用了,获取客户端连接时的ip就是通过反射来获取对应字段的属性来获取的。
WebSocket 获取客户端的IP_websocket获取当前客户端ip-CSDN博客
这种方式暂未试验,通过拦截器获取ip
WebSocket 实现长连接及通过WebSocket获取客户端IP_websocket 获取客户端ip-CSDN博客
4.其他属性
根据自己的实际情况来定义,像是存储的ip,端口等。
或者是存储websocket类,通过可区分的唯一id来存储该类,可以获取该类的所有属性

3.2方法
分别是连接时调用方法,发送信息,发生异常时,关闭连接时。
1.连接方法
/**
* 连接成功调用的方法
*/
@OnOpen
public void onOpen(Session session, @PathParam("sid") String param) {
try {
String newIp = WebsocketUtil.getRemoteAddress(session);
String clientIp = newIp.substring(newIp.lastIndexOf("=") + 2,newIp.length() - 1);
this.session = session;
clients.put(clientIp, this);
//在线数加1
addOnlineCount();
if ("0".equals(param)) {
this.webIp = clientIp;
this.ip = clientIp;
sid = "web";
sendMessage("WebSocket连接成功!");
log.info("有新窗口开始监听,客户端为web端,ip为:" + webIp + "当前在线数量为:" + getOnlineCount());
} else if ("1".equals(param)) {
String hostname = session.getQueryString().substring(session.getQueryString().indexOf("=") + 1);
this.hostname = hostname;
//存放处理后pc端ip
list.add(clientIp.substring(0, clientIp.indexOf(":")) + "-" + hostname);
this.pcIp = clientIp;
this.ip = clientIp;
sid = "pc";
sendMessage("WebSocket连接成功!");
log.info("有新窗口开始监听,客户端为pc端,ip为:" + pcIp + "当前在线数量为:" + getOnlineCount());
}
this.port = ip.substring(ip.indexOf(":") + 1);
} catch (IOException e) {
log.error("ip为:" + ip + "客户端连接异常");
}
}
2.关闭方法
/**
* 连接关闭调用方法
*/
@OnClose
public void onClose() {
webSocketSet.remove(this); //从set中删除;
if (pcIp != null && !pcIp.isEmpty()) {
list.remove(pcIp.substring(0, pcIp.indexOf(":")) + "-" + hostname);
}
clients.remove(ip);
subOnlineCount(); //在线数减1
if ("web".equals(sid)) {
log.info("web客户端ip为:" + ip + "的WebSocket连接关闭!当前在线数量为:" + getOnlineCount());
} else if ("pc".equals(sid)) {
log.info("pc客户端ip为:" + ip + "的WebSocket连接关闭!当前在线数量为:" + getOnlineCount());
}
}
3.消息转发方法
this.session.getBasicRemote().sendText(message);
StringUtils.ordinalIndexOf 是 Apache Commons Lang 库中的一个静态方法,用于查找字符串中某个子字符串(目标字符串)出现的第 n 次位置的索引。它的作用是返回目标字符串第 n 次出现在源字符串中的索引位置。
Arrays.asList() 是将数组转换为 List 的便捷方法,可以在数组和 List 之间方便地进行数据操作和转换。需要注意的是返回的
List 是固定大小的,不能进行结构性修改操作,同时修改 List 的元素会直接反映到原数组中。
/**
* 接收客户端发送的信息,并进行转发
*/
@OnMessage
public void onMessage(String message, Session session) throws IOException {
if ("0".equals(message)) {
log.debug("pc客户端:" + ip + "连接中");
session.getBasicRemote().sendText("1");
} else if ("开始上传日志".equals(message)) {
System.out.println(message);
} else if (message.contains("Y") || message.contains("N")){
String returnMessage = message.substring(0,StringUtils.ordinalIndexOf(message, ",", 2));
List<String> ipList = Arrays.asList(message.substring(
StringUtils.ordinalIndexOf(message, ",", 2) + 1).split(","));
for (Map.Entry<String, WebSocketServer> entry : clients.entrySet()) {
for (int i = 0; i < ipList.size(); i++) {
String forwardIp = ipList.get(i);
if (entry.getKey().contains(forwardIp.substring(0,forwardIp.indexOf("-"))) && "pc".equals(entry.getValue().sid)) {
forwardIp = entry.getKey();
s_cIps.put(forwardIp, ip);
log.info("收到来自web客户端:" + ip + "的信息:" + message);
clients.get(forwardIp).session.getBasicRemote().sendText(returnMessage);
log.info("发送给pc客户端:" + forwardIp + "的信息为:" + returnMessage);
session.getBasicRemote().sendText("发送成功");
}else if (ip(forwardIp)){
log.info("发送失败!客户端ip为:" + forwardIp + "可能断开连接了,请刷新ip");
session.getBasicRemote().sendText("发送失败!客户端ip为:" + forwardIp + "可能断开连接了,请刷新ip");
}
}
}
} else {
String forwardIp = message.substring(0, message.indexOf(","));
String returnMessage = message.replace(forwardIp + ",", "");
forwardIp = forwardIp.substring(0,forwardIp.indexOf("-"));
for (Map.Entry<String, WebSocketServer> entry : clients.entrySet()) {
if (entry.getKey().contains(forwardIp) && "pc".equals(entry.getValue().sid)) {
forwardIp = entry.getKey();
}
}
try {
s_cIps.put(forwardIp, ip);
log.info("收到来自web客户端:" + ip + "的信息:" + returnMessage);
clients.get(forwardIp).session.getBasicRemote().sendText(returnMessage);
log.info("发送给pc客户端:" + forwardIp + "的信息为:" + returnMessage);
session.getBasicRemote().sendText("发送成功");
} catch (NullPointerException e) {
log.error("发送失败!客户端ip为:" + forwardIp + "可能断开连接了");
session.getBasicRemote().sendText("发送失败!客户端ip为:" + forwardIp + "可能断开连接了");
}
}
}
4.连接异常方法
/**
* 连接出错
*
* @param session
* @param throwable
*/
@OnError
public void onError(Session session, Throwable throwable) {
log.error("ip为:" + ip + "的" + sid + "客户端发生错误 异常信息:" + throwable.getMessage());
}
5.其他方法
通过ip回传消息
Map.Entry 是 Java 中表示键值对的接口,它是嵌套在 Map 接口中的静态接口。Map.Entry 表示了键值对的条目,包括键和对应的值。通过 Map.Entry,可以方便地遍历和操作 Map 中的键值对数据。
entrySet() 返回值的特点:
1.类型: entrySet() 方法返回一个实现了 Set> 接口的集合,其中 K 是键的类型,V 是值的类型。
2.可变性: 通过返回的 Set 集合,可以直接修改 Map 中的内容。例如,使用 Map.Entry 的 setValue 方法可以修改对应的值。
3.实时反映: entrySet() 返回的集合视图是实时反映 Map 中的变化的,如果 Map 被修改(添加、删除、更新条目),则 entrySet 视图也会随之变化。
/**
* 回传消息
*
* @param ip
* @param msg
* @return
*/
private static ConcurrentHashMap<String, WebSocket> clients = new ConcurrentHashMap<>();
public static String sendMsg(String ip, String msg) {
try {
for (Map.Entry<String, WebSocketServer> entry : clients.entrySet()) {
if (entry.getKey().contains(ip) && "pc".equals(entry.getValue().sid)) {
clients.get(s_cIps.get(entry.getKey())).session.getBasicRemote().sendText(msg);
}
}
} catch (Exception e) {
return "消息发送失败";
}
return "消息发送成功";
}
更多推荐




所有评论(0)