全双工通信: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博客

Java 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 "消息发送成功";
}

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐