feat(websocket): 实现基于Redis的消息队列WebSocket通信功能

- 添加Redis消息队列配置类RedisMessageConfig,实现消息监听容器
- 重构KitWsMessage消息结构,将数据字段改为内部Data类并添加getter/setter
- 集成RedisTemplate实现WebSocket消息发布功能,支持主动推送和实时聊天
- 实现WsMessageSubscriber消息订阅者,处理来自Redis的消息队列数据
- 添加WebSocket业务代码常量接口KitWsBizCodeConst
- 更新前端Demo组件适配新的消息格式和业务代码参数
- 优化WebSocket服务器端消息处理逻辑和错误处理机制
This commit is contained in:
2026-02-25 13:18:01 +08:00
parent ae0c0482cf
commit d74bcd5337
13 changed files with 454 additions and 54 deletions
@@ -63,19 +63,19 @@ export default {
console.info('=============== KitWsClient:event', event) console.info('=============== KitWsClient:event', event)
}, function(msgEvent) { }, function(msgEvent) {
console.info('=============== KitWsClient:msgEvent', msgEvent) console.info('=============== KitWsClient:msgEvent', msgEvent)
that.responseTxt = that.responseTxt + msgEvent.data + '' that.responseTxt = that.responseTxt + msgEvent.data + '\n'
}) })
}, },
onSendClick() { onSendClick() {
if (this.client) { const that = this
this.client.sendData({ if (that.client) {
msgType: 'sc', that.client.sendData({
userIds: ['123456'], bizCode: '101',
content: this.sendVal data: {
}) fromUsername: that.user.username,
this.client.sendData({ toUsername: null,
msgType: 'gc', message: that.sendVal
content: this.sendVal }
}) })
} }
}, },
@@ -0,0 +1,60 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.redis;
import cn.odboy.framework.websocket.WsMessageSubscriber;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.listener.ChannelTopic;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.data.redis.listener.adapter.MessageListenerAdapter;
/**
* Redis消息队列配置
*
* @author odboy
*/
@Configuration
public class RedisMessageConfig {
/**
* 消息监听容器
*
* @param connectionFactory /
* @param listenerAdapter /
* @return /
*/
@Bean
public RedisMessageListenerContainer container(RedisConnectionFactory connectionFactory, MessageListenerAdapter listenerAdapter) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(connectionFactory);
// 指定监听的主题
container.addMessageListener(listenerAdapter, new ChannelTopic("WebSocketMessage"));
return container;
}
/**
* 将订阅者与消息处理方法绑定
*
* @param subscriber /
* @return /
*/
@Bean
public MessageListenerAdapter listenerAdapter(WsMessageSubscriber subscriber) {
return new MessageListenerAdapter(subscriber, "onMessage");
}
}
@@ -0,0 +1,27 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.websocket;
/**
* WebSocket业务代码常量
*
* @author odboy
*/
public interface KitWsBizCodeConst {
String AutomaticPush = "100"; // 主动推送
String IM = "101"; // 实时聊天
}
@@ -0,0 +1,55 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.websocket;
import cn.odboy.framework.websocket.context.KitWsMessage;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class KitWsMessagePublisher {
@Autowired
private RedisTemplate<Object, Object> redisTemplate;
private KitWsMessage innerBuild(String bizCode, String fromUser, String toUser, String message) {
KitWsMessage wsMessage = new KitWsMessage();
wsMessage.setBizCode(bizCode);
KitWsMessage.Data data = new KitWsMessage.Data();
data.setFormUser(fromUser);
data.setToUser(toUser);
data.setMessage(message);
wsMessage.setData(data);
return wsMessage;
}
public void automaticPush(String fromUser, String toUser, String message) {
KitWsMessage wsMessage = this.innerBuild(KitWsBizCodeConst.AutomaticPush, fromUser, toUser, message);
redisTemplate.convertAndSend("WebSocketMessage", wsMessage);
log.info("发送AutomaticPush消息");
}
public void imPush(String fromUser, String toUser, String message) {
KitWsMessage wsMessage = this.innerBuild(KitWsBizCodeConst.IM, fromUser, toUser, message);
redisTemplate.convertAndSend("WebSocketMessage", wsMessage);
log.info("发送IM消息");
}
}
@@ -0,0 +1,51 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.websocket;
import cn.odboy.framework.websocket.context.KitWsMessage;
import cn.odboy.framework.websocket.context.KitWsServer;
import com.alibaba.fastjson2.JSON;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
/**
* Ws消息订阅
*
* @author odboy
*/
@Slf4j
@Component
public class WsMessageSubscriber {
@Autowired
private KitWsServer kitWsServer;
// /**
// * 无效写法
// */
// public void onMessage(KitWsMessage message) {
// log.info("收到来自redis消息队列的数据, {}", JSON.toJSONString(message));
// KitWsMessage.Data data = message.getData();
// kitWsServer.sendMessage(data.getMessage(), data.getToUsername());
// }
public void onMessage(String message, String channel) {
log.info("收到来自redis消息队列的message={}, channel={}", JSON.toJSONString(message), channel);
KitWsMessage wsMessage = JSON.parseObject(message, KitWsMessage.class);
kitWsServer.sendMessage(wsMessage.getData().getMessage(), wsMessage.getData().getToUser());
}
}
@@ -25,10 +25,14 @@ import lombok.Setter;
public class KitWsMessage extends KitObject { public class KitWsMessage extends KitObject {
private String bizCode; private String bizCode;
private String data; private Data data;
public KitWsMessage(String bizCode, String data) { @Getter
this.bizCode = bizCode; @Setter
this.data = data; public static class Data extends KitObject {
private String formUser;
private String toUser;
private String message;
} }
} }
@@ -13,23 +13,22 @@
* See the License for the specific language governing permissions and * See the License for the specific language governing permissions and
* limitations under the License. * limitations under the License.
*/ */
package cn.odboy.framework.websocket.context; package cn.odboy.framework.websocket.context;
import cn.hutool.core.util.StrUtil;
import cn.odboy.framework.context.KitSpringBeanHolder;
import cn.odboy.framework.websocket.KitWsMessagePublisher;
import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSON;
import jakarta.websocket.OnClose; import jakarta.websocket.*;
import jakarta.websocket.OnError;
import jakarta.websocket.OnMessage;
import jakarta.websocket.OnOpen;
import jakarta.websocket.Session;
import jakarta.websocket.server.PathParam; import jakarta.websocket.server.PathParam;
import jakarta.websocket.server.ServerEndpoint; import jakarta.websocket.server.ServerEndpoint;
import java.io.IOException;
import java.util.Objects;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.io.IOException;
import java.util.Objects;
@Slf4j @Slf4j
@Component @Component
@ServerEndpoint("/websocket/{sid}") @ServerEndpoint("/websocket/{sid}")
@@ -51,7 +50,6 @@ public class KitWsServer {
@OnOpen @OnOpen
public void onOpen(Session session, @PathParam("sid") String sid) { public void onOpen(Session session, @PathParam("sid") String sid) {
this.session = session; this.session = session;
// {username}#{bizCode}#{contextParams}
this.sid = sid; this.sid = sid;
KitWsClientManager.addClient(sid, this); KitWsClientManager.addClient(sid, this);
} }
@@ -64,10 +62,12 @@ public class KitWsServer {
@OnMessage @OnMessage
public void onMessage(String message, Session session) throws IOException { public void onMessage(String message, Session session) throws IOException {
KitWsMessage wsMessage = JSON.parseObject(message, KitWsMessage.class); KitWsMessage wsMessage = JSON.parseObject(message, KitWsMessage.class);
String bizCode = wsMessage.getBizCode(); // String bizCode = wsMessage.getBizCode();
Object data = wsMessage.getData(); // Object data = wsMessage.getData();
log.info("收到来 sid={} 的信息: message={}, bizCode={}, data={}", sid, message, bizCode, JSON.toJSONString(data)); log.info("收到来 sid={} 的信息: message={}", sid, message);
session.getBasicRemote().sendText(message); // session.getBasicRemote().sendText(message);
KitWsMessagePublisher messagePublisher = KitSpringBeanHolder.getBean(KitWsMessagePublisher.class);
messagePublisher.imPush(wsMessage.getData().getFormUser(), wsMessage.getData().getToUser(), wsMessage.getData().getMessage());
} }
/** /**
@@ -84,7 +84,7 @@ public class KitWsServer {
@OnError @OnError
public void onError(Session session, Throwable error) { public void onError(Session session, Throwable error) {
// log.error("WebSocket sid={} 发生错误", this.sid, error); log.error("WebSocket sid={} 发生错误", this.sid, error);
try { try {
KitWsClientManager.removeClient(this.sid); KitWsClientManager.removeClient(this.sid);
} catch (Exception e) { } catch (Exception e) {
@@ -102,26 +102,26 @@ public class KitWsServer {
/** /**
* 精准推送消息 * 精准推送消息
*/ */
public void sendMessage(KitWsMessage message, @PathParam("sid") String sid) throws IOException { public void sendMessage(String message, @PathParam("sid") String sid) {
String body = JSON.toJSONString(message); log.info("推送消息到{}, 推送内容:{}", sid, message);
// log.info("推送消息到{}, 推送内容:{}", sid, message);
try { try {
if (sid == null) { if (StrUtil.isBlank(sid)) {
sendToAll(body); sendToAll(message);
} else { } else {
KitWsServer wsClient = KitWsClientManager.getClientBySid(sid); KitWsServer wsClient = KitWsClientManager.getClientBySid(sid);
if (wsClient != null) { if (wsClient != null) {
wsClient.innerSendMessage(body); wsClient.innerSendMessage(message);
} }
} }
} catch (Exception ignored) { } catch (IOException e) {
log.warn("推送消息到客户端失败", e);
} }
} }
/** /**
* 群发消息 * 群发消息
*/ */
public void sendToAll(String message) { private void sendToAll(String message) {
for (KitWsServer item : KitWsClientManager.getAllClient()) { for (KitWsServer item : KitWsClientManager.getAllClient()) {
try { try {
item.innerSendMessage(message); item.innerSendMessage(message);
@@ -0,0 +1,60 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.redis;
import cn.odboy.framework.websocket.WsMessageSubscriber;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.listener.ChannelTopic;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.data.redis.listener.adapter.MessageListenerAdapter;
/**
* Redis消息队列配置
*
* @author odboy
*/
@Configuration
public class RedisMessageConfig {
/**
* 消息监听容器
*
* @param connectionFactory /
* @param listenerAdapter /
* @return /
*/
@Bean
public RedisMessageListenerContainer container(RedisConnectionFactory connectionFactory, MessageListenerAdapter listenerAdapter) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(connectionFactory);
// 指定监听的主题
container.addMessageListener(listenerAdapter, new ChannelTopic("WebSocketMessage"));
return container;
}
/**
* 将订阅者与消息处理方法绑定
*
* @param subscriber /
* @return /
*/
@Bean
public MessageListenerAdapter listenerAdapter(WsMessageSubscriber subscriber) {
return new MessageListenerAdapter(subscriber, "onMessage");
}
}
@@ -0,0 +1,27 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.websocket;
/**
* WebSocket业务代码常量
*
* @author odboy
*/
public interface KitWsBizCodeConst {
String AutomaticPush = "100"; // 主动推送
String IM = "101"; // 实时聊天
}
@@ -0,0 +1,55 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.websocket;
import cn.odboy.framework.websocket.context.KitWsMessage;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class KitWsMessagePublisher {
@Autowired
private RedisTemplate<Object, Object> redisTemplate;
private KitWsMessage innerBuild(String bizCode, String fromUser, String toUser, String message) {
KitWsMessage wsMessage = new KitWsMessage();
wsMessage.setBizCode(bizCode);
KitWsMessage.Data data = new KitWsMessage.Data();
data.setFormUser(fromUser);
data.setToUser(toUser);
data.setMessage(message);
wsMessage.setData(data);
return wsMessage;
}
public void automaticPush(String fromUser, String toUser, String message) {
KitWsMessage wsMessage = this.innerBuild(KitWsBizCodeConst.AutomaticPush, fromUser, toUser, message);
redisTemplate.convertAndSend("WebSocketMessage", wsMessage);
log.info("发送AutomaticPush消息");
}
public void imPush(String fromUser, String toUser, String message) {
KitWsMessage wsMessage = this.innerBuild(KitWsBizCodeConst.IM, fromUser, toUser, message);
redisTemplate.convertAndSend("WebSocketMessage", wsMessage);
log.info("发送IM消息");
}
}
@@ -0,0 +1,51 @@
/*
* Copyright 2021-2026 Odboy
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package cn.odboy.framework.websocket;
import cn.odboy.framework.websocket.context.KitWsMessage;
import cn.odboy.framework.websocket.context.KitWsServer;
import com.alibaba.fastjson2.JSON;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
/**
* Ws消息订阅
*
* @author odboy
*/
@Slf4j
@Component
public class WsMessageSubscriber {
@Autowired
private KitWsServer kitWsServer;
// /**
// * 无效写法
// */
// public void onMessage(KitWsMessage message) {
// log.info("收到来自redis消息队列的数据, {}", JSON.toJSONString(message));
// KitWsMessage.Data data = message.getData();
// kitWsServer.sendMessage(data.getMessage(), data.getToUsername());
// }
public void onMessage(String message, String channel) {
log.info("收到来自redis消息队列的message={}, channel={}", JSON.toJSONString(message), channel);
KitWsMessage wsMessage = JSON.parseObject(message, KitWsMessage.class);
kitWsServer.sendMessage(wsMessage.getData().getMessage(), wsMessage.getData().getToUser());
}
}
@@ -25,10 +25,14 @@ import lombok.Setter;
public class KitWsMessage extends KitObject { public class KitWsMessage extends KitObject {
private String bizCode; private String bizCode;
private String data; private Data data;
public KitWsMessage(String bizCode, String data) { @Getter
this.bizCode = bizCode; @Setter
this.data = data; public static class Data extends KitObject {
private String formUser;
private String toUser;
private String message;
} }
} }
@@ -15,7 +15,11 @@
*/ */
package cn.odboy.framework.websocket.context; package cn.odboy.framework.websocket.context;
import cn.hutool.core.util.StrUtil;
import cn.odboy.framework.context.KitSpringBeanHolder;
import cn.odboy.framework.websocket.KitWsMessagePublisher;
import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSON;
import java.io.IOException; import java.io.IOException;
import java.util.Objects; import java.util.Objects;
import javax.websocket.OnClose; import javax.websocket.OnClose;
@@ -25,6 +29,7 @@ import javax.websocket.OnOpen;
import javax.websocket.Session; import javax.websocket.Session;
import javax.websocket.server.PathParam; import javax.websocket.server.PathParam;
import javax.websocket.server.ServerEndpoint; import javax.websocket.server.ServerEndpoint;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@@ -50,7 +55,6 @@ public class KitWsServer {
@OnOpen @OnOpen
public void onOpen(Session session, @PathParam("sid") String sid) { public void onOpen(Session session, @PathParam("sid") String sid) {
this.session = session; this.session = session;
// {username}#{bizCode}#{contextParams}
this.sid = sid; this.sid = sid;
KitWsClientManager.addClient(sid, this); KitWsClientManager.addClient(sid, this);
} }
@@ -63,10 +67,12 @@ public class KitWsServer {
@OnMessage @OnMessage
public void onMessage(String message, Session session) throws IOException { public void onMessage(String message, Session session) throws IOException {
KitWsMessage wsMessage = JSON.parseObject(message, KitWsMessage.class); KitWsMessage wsMessage = JSON.parseObject(message, KitWsMessage.class);
String bizCode = wsMessage.getBizCode(); // String bizCode = wsMessage.getBizCode();
Object data = wsMessage.getData(); // Object data = wsMessage.getData();
log.info("收到来 sid={} 的信息: message={}, bizCode={}, data={}", sid, message, bizCode, JSON.toJSONString(data)); log.info("收到来 sid={} 的信息: message={}", sid, message);
session.getBasicRemote().sendText(message); // session.getBasicRemote().sendText(message);
KitWsMessagePublisher messagePublisher = KitSpringBeanHolder.getBean(KitWsMessagePublisher.class);
messagePublisher.imPush(wsMessage.getData().getFormUser(), wsMessage.getData().getToUser(), wsMessage.getData().getMessage());
} }
/** /**
@@ -83,7 +89,7 @@ public class KitWsServer {
@OnError @OnError
public void onError(Session session, Throwable error) { public void onError(Session session, Throwable error) {
// log.error("WebSocket sid={} 发生错误", this.sid, error); log.error("WebSocket sid={} 发生错误", this.sid, error);
try { try {
KitWsClientManager.removeClient(this.sid); KitWsClientManager.removeClient(this.sid);
} catch (Exception e) { } catch (Exception e) {
@@ -101,26 +107,26 @@ public class KitWsServer {
/** /**
* 精准推送消息 * 精准推送消息
*/ */
public void sendMessage(KitWsMessage message, @PathParam("sid") String sid) throws IOException { public void sendMessage(String message, @PathParam("sid") String sid) {
String body = JSON.toJSONString(message); log.info("推送消息到{}, 推送内容:{}", sid, message);
// log.info("推送消息到{}, 推送内容:{}", sid, message);
try { try {
if (sid == null) { if (StrUtil.isBlank(sid)) {
sendToAll(body); sendToAll(message);
} else { } else {
KitWsServer wsClient = KitWsClientManager.getClientBySid(sid); KitWsServer wsClient = KitWsClientManager.getClientBySid(sid);
if (wsClient != null) { if (wsClient != null) {
wsClient.innerSendMessage(body); wsClient.innerSendMessage(message);
} }
} }
} catch (Exception ignored) { } catch (IOException e) {
log.warn("推送消息到客户端失败", e);
} }
} }
/** /**
* 群发消息 * 群发消息
*/ */
public void sendToAll(String message) { private void sendToAll(String message) {
for (KitWsServer item : KitWsClientManager.getAllClient()) { for (KitWsServer item : KitWsClientManager.getAllClient()) {
try { try {
item.innerSendMessage(message); item.innerSendMessage(message);