前言

        版本注意,使用的是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;


}

Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐