Pig4cloud学习版使用websocket的过程
前言
版本注意,使用的是pig4cloud学习版3.8.2版本。由于学习版本文档显示的内容不能很全,导致我使用整合的websocket功能过程出现问题,不能按照网站的示例进行,反而是问了半天ai和看了半天源码才跑通代码。我是想在新建立的微服务中,进行前后端通信功能,比如后端预警数据库表出现新数据,则通过websocket传到前端,前端收到消息在页面上提醒。下面是结果展示,前端页面出现预警消息。
接下来,说一下pig4cloud前端的部分,在src/components里面有他封装好的websocket,可以通过引用直接使用这里的写好的ws等内容。
// ws地址
let host = window.location.host;
// baseURL
let baseURL = import.meta.env.VITE_API_URL;
let wsUri = `ws://${host}${baseURL}${other.adaptationUrl(props.uri)}?access_token=${token.value}&TENANT-ID=${tenant.value}`;
// 建立连接
state.webSocket = new WebSocket(wsUri);
这里是里面的建立ws前端的代码,你可以看到,他已经写好了wsUri,其中里面的host,baseURL(这里的uri是上面通过模块添加的,主要是和后端保持一致,已经token和tenant,上面也获取,这些先不需要考虑)。这里你可以打印wsUri是什么,通过看env,env.development,vite.config.ts文件,可以对应的服务地址找到,发现是admin服务地址。如果你后端是在admin服务里面写websocket没问题。但是我是新建立的waterresource服务,服务地址都不一样,所以就是不能用这里封装好的websocket.vue,除非你去改相关内容。在env里面记得修改
# 是否开启websocket 消息接受,
VITE_WEBSOCKET_ENABLE = ture,改为ture就行。



其中pig4cloud整合websocket, 核心功能如下:
| 功能 | 说明 |
|---|---|
| WebSocket 服务端 | 自动配置 WebSocket 服务端,支持 @ServerEndpoint 注解开发 |
| 分布式消息广播 | 基于 Redis PUB/SUB 实现多节点间的消息广播(适用于集群部署) |
| 会话管理 | 提供 WebSocketSessionHolder 工具类管理活跃会话 |
| 消息分发器 | 通过 RedisMessageDistributor 实现消息的定向推送或广播 |
| 自动健康检查 | 集成 Spring Boot Actuator,暴露 /actuator/websocket 端点 |
| 自定义协议支持 | 支持 STOMP、自定义二进制协议等(需配合额外配置) |
你可以使用application等配置文件,快速配置构建websocket服务。至于里面写好的类,你可以看一看他封装好的源码,至于使用,我下面会展示一种会话发送,不是很全(也可以问ai)(但是可能很便秘),这里封装的类你可以进行覆盖,如果你对sping和websocket很了解的话,可以直接写符合你项目要求的,如果只是发送消息去前端,直接用他封装好的就好了,清清爽爽。
前端
<el-alert v-if="showWarning" :title="`${warningData.station} - ${warningData.content}`"
:type="warningData.level.includes('一级') ? 'error' : 'warning'" center show-icon :closable="true"
class="global-alert" @close="showWarning = false" />
import { ElNotification } from 'element-plus';
import { Session } from '/@/utils/storage';
let ws = null;
const maxRetries = 3; // 最大重试次数
let retryCount = 0;
const showWarning = ref(false);
const warningData = ref({});
function connectWebSocket() {
const token = Session.getToken();
const tenant = Session.getTenant();
const wsUri = `ws://localhost:7002/waterresource/ws/info?access_token=${token}&TENANT-ID=${tenant}`;
ws = new WebSocket(wsUri);
ws.onopen = () => {
console.log("WebSocket连接成功");
retryCount = 0; // 重置重试计数器
};
ws.onmessage = (event) => {
// console.log("收到消息:", event.data);
try {
const data = JSON.parse(event.data);
warningData.value = data;
showWarning.value = true;
ElNotification.warning({
title: '告警通知',
message: data.content || event.data,
type: 'warning',
duration: 0 // 消息不自动关闭
});
} catch (e) {
console.error("解析消息失败:", e);
}
};
ws.onerror = (e) => {
console.error("WebSocket错误:", e);
ws.close(); // 关闭错误连接
};
ws.onclose = (e) => {
console.log("WebSocket连接关闭:", e);
if (retryCount < maxRetries) {
retryCount++;
console.log(`尝试第 ${retryCount} 次重连...`);
setTimeout(connectWebSocket, 3000); // 3秒后重连
}
};
}
.global-alert {
position: absolute;
/* 使用 fixed 定位确保始终居中 */
top: 10%;
left: 50%;
/* 水平居中 */
transform: translateX(-50%);
/* 向左偏移自身宽度的一半 */
width: 80%;
/* 宽度设置为 80% */
max-width: 800px;
/* 最大宽度限制 */
z-index: 9999;
border-radius: 4px;
box-shadow: 0 2px 12px rgba(0, 0, 0, 0.15);
}
这里写了相关对应的代码部分,自己提取想要的。这里是硬编码,我直接写了
const wsUri = `ws://localhost:7002/waterresource/ws/info?access_token=${token}&TENANT-ID=${tenant}`;
这里的waterresource/ws/info的地址要和后端websocket配置里面的path对应一致,以及我的waterreource服务端口是7002也是要对应的,别发送的admin服务去了,到时候就连接不到了。这里的token和tennant就和他封装好示例里的一样就好了,可以参考,我这里只是将前面确定的部分换了一下。以及写了一些ElNotification和el-alert展示,你先确定websocket可以连接到后端在慢慢调试其他部分。这里就是前端部分了。
后端
1.添加依赖
<dependency>
<groupId>com.pig4cloud.plugin</groupId>
<artifactId>websocket-spring-boot-starter</artifactId>
<version>3.0.0</version>
</dependency>
2.配置application.yml
注意这类的path要和我前面说的path与前端对应,你也可以改,只要一样就行。比如/abc。注意如果开启redis,你也要加redis的相关内容。

3.WebSocketConfig配置类
这里就是从session中提取想要的内容,以及生成sessionkeys,这里是根据你的需要去获取内容,因为后面他们发送到指定地方是通过sessionkey识别的,你可以轮询。但是你得确保服务端可以找到会话集合,这样才可以去发消息。
package com.pig4cloud.pig.waterresource.config;
import com.pig4cloud.plugin.websocket.custom.SecuritySessionKeyGenerator;
import com.pig4cloud.plugin.websocket.holder.SessionKeyGenerator;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.codec.digest.DigestUtils;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.WebSocketSession;
import java.net.URI;
import java.util.*;
import java.util.stream.Collectors;
@Slf4j
@Configuration
public class WebSocketConfig {
@Bean
public SessionKeyGenerator sessionKeyGenerator() {
return new SecuritySessionKeyGenerator() {
@Override
public Object sessionKey(WebSocketSession session) {
try {
Map<String, String> params = parseQueryParams(session.getUri());
String tenantId = params.get("TENANT-ID");
String token = params.get("access_token");
if (tenantId == null || token == null) {
throw new IllegalArgumentException("参数缺失");
}
// 生成 sessionKey 并存储
String sessionKey = String.format("%s_%s", tenantId,
DigestUtils.md5Hex(token).substring(0, 8));
// 关键点:将 sessionKey 存入 session 的 attributes
session.getAttributes().put("sessionKey", sessionKey);
session.getAttributes().put("tenantId", tenantId);
session.getAttributes().put("token", token);
log.info("生成 sessionKey: {}", sessionKey);
return sessionKey;
} catch (Exception e) {
log.error("生成 sessionKey 失败: {}", e.getMessage());
return "fallback_" + UUID.randomUUID();
}
}
};
}
private Map<String, String> parseQueryParams(URI uri) {
if (uri == null || uri.getQuery() == null) {
return Collections.emptyMap();
}
return Arrays.stream(uri.getQuery().split("&"))
.map(param -> param.split("="))
.filter(arr -> arr.length == 2)
.collect(Collectors.toMap(
arr -> arr[0], // 保持原始大小写
arr -> arr[1],
(v1, v2) -> v1
));
}
}
4.WarningRecordRabbitListener监听类
注意,我是使用rabbitmq等方法去监听数据库变化,并且发送到这里,你另外配置你想要的技术即可,我这里只展示websocket相关内容了,你可以直接通过某个接口调用这里的功能也行。以及我这里可以获取到WebSocketConfig存入的sessionkeys等内容,那也就可以发送到指定地方了。以及我这里的tenantId是硬编码,你可以在对应的dto中,或者数据库表中获取,我这里就直接上了(快绷不住了)(至于为什么知道是,前端页面看的网络请求地址)
package com.pig4cloud.pig.waterresource.listener;
import com.pig4cloud.pig.waterresource.dto.WarningMessageDTO;
import com.pig4cloud.plugin.websocket.distribute.MessageDO;
import com.pig4cloud.plugin.websocket.distribute.RedisMessageDistributor;
import com.pig4cloud.plugin.websocket.holder.WebSocketSessionHolder;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
import org.springframework.web.socket.WebSocketSession;
import java.util.*;
import java.util.stream.Collectors;
@Slf4j
@Component
@RequiredArgsConstructor
public class WarningRecordRabbitListener {
private final RedisMessageDistributor messageDistributor;
@RabbitListener(queues = "warning.record.queue")
public void handleNewWarningRecord(WarningMessageDTO message) {
// 打印所有会话信息(调试用)
WebSocketSessionHolder.getSessions().forEach(session -> {
log.info("活跃会话详情 >>> ID={}, tenantId={}, sessionKey={}, attributes={}",
session.getId(),
session.getAttributes().get("tenantId"),
session.getAttributes().get("sessionKey"),
session.getAttributes()
);
});
// 获取目标会话的 sessionKey
List<Object> targetSessionKeys = WebSocketSessionHolder.getSessions().stream()
.filter(session -> {
Object tenantId = session.getAttributes().get("tenantId");
return tenantId != null && tenantId.toString().equals("1"); // 假设租户ID为1
})
.map(session -> session.getAttributes().get("sessionKey")) // 使用 sessionKey 推送
.filter(Objects::nonNull)
.collect(Collectors.toList());
if (targetSessionKeys.isEmpty()) {
log.warn("没有匹配的活跃会话,消息丢弃");
return;
}
// 推送消息
MessageDO messageDO = new MessageDO()
.setNeedBroadcast(false)
.setSessionKeys(targetSessionKeys)
.setMessageText(convertToMessage(message));
messageDistributor.distribute(messageDO);
log.info("消息已推送至会话: {}", targetSessionKeys);
}
/**
* 从消息中获取关联的租户ID(根据实际业务逻辑调整)
*/
private String getTenantIdFromMessage(WarningMessageDTO message) {
// 方案1:如果消息中有设备ID等信息,通过设备ID查询租户ID
// 假设存在设备服务可以根据设备名称查询租户
// return deviceService.getTenantIdByStationName(message.getStationName());
// 方案2:如果所有消息都属于同一个租户(测试环境)
return "1"; // 假设默认租户ID为1
// 方案3:如果消息中包含租户相关信息(如通过stationName解析)
/*
String stationName = message.getStationName();
// 假设stationName格式为"租户ID_站点名称"
if (stationName != null && stationName.contains("_")) {
return stationName.split("_")[0];
}
return "1"; // 默认租户
*/
}
/**
* 消息内容转换
*/
private String convertToMessage(WarningMessageDTO dto) {
return String.format("{\"id\":%d,\"station\":\"%s\",\"content\":\"%s\",\"level\":\"%s\"}",
dto.getId(),
dto.getStationName(),
dto.getContent(),
dto.getLevel());
}
}
这里是我的dto类,主要是返回出现问题的地点名字,报警内容,报警等级,报警时间。没有租户id等,还是不够完善,但是发送一下没问题。
package com.pig4cloud.pig.waterresource.dto;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
@Data
@NoArgsConstructor
@AllArgsConstructor
@Builder
public class WarningMessageDTO {
private Integer id;
private String stationName;
private String content;
private String level;
private LocalDateTime createTime;
}
更多推荐


所有评论(0)