springboot集成集群模式websocket服务
引言
springboot项目集成websocket服务,并且使用redis发布订阅方案解决websocket集群模式下session共享问题
步骤
1.导包
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.hui</groupId>
<artifactId>websocket_demo</artifactId>
<version>1.0.0-SNAPSHOT</version>
<properties>
<java.version>8</java.version>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
<spring-boot.version>2.3.7.RELEASE</spring-boot.version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-dependencies</artifactId>
<version>${spring-boot.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
<exclusions>
<exclusion>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-tomcat</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-undertow</artifactId>
<version>3.0.1</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-pool2</artifactId>
<version>2.11.1</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>cn.hutool</groupId>
<artifactId>hutool-all</artifactId>
<version>5.8.5</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
</plugin>
</plugins>
</build>
<repositories>
<repository>
<id>spring-milestones</id>
<name>Spring Milestones</name>
<url>https://repo.spring.io/milestone</url>
<snapshots>
<enabled>false</enabled>
</snapshots>
</repository>
<repository>
<id>spring-snapshots</id>
<name>Spring Snapshots</name>
<url>https://repo.spring.io/snapshot</url>
<releases>
<enabled>false</enabled>
</releases>
</repository>
</repositories>
<pluginRepositories>
<pluginRepository>
<id>spring-milestones</id>
<name>Spring Milestones</name>
<url>https://repo.spring.io/milestone</url>
<snapshots>
<enabled>false</enabled>
</snapshots>
</pluginRepository>
<pluginRepository>
<id>spring-snapshots</id>
<name>Spring Snapshots</name>
<url>https://repo.spring.io/snapshot</url>
<releases>
<enabled>false</enabled>
</releases>
</pluginRepository>
</pluginRepositories>
</project>
2. 配置文件
- yml文件
server:
port: 9000
spring:
redis:
host: 10.99.11.40
port: 6379
password: HeyGears@2022
database: 15
timeout: 9000
lettuce:
pool:
max-active: 100
min-idle: 5
max-wait: -1
max-idle: 10
3. 代码
MessageDto :
@Data
@AllArgsConstructor
@NoArgsConstructor
public class MessageDto implements Serializable {
private static final long serialVersionUID = -4291728346293647762L;
private String toUid;
private String sendUid;
private String msg;
}
RedisConfig :
@Configuration
@EnableCaching
public class RedisConfig extends CachingConfigurerSupport {
@Bean
public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory connectionFactory) {
RedisTemplate<String, Object> template = new RedisTemplate<>();
template.setConnectionFactory(connectionFactory);
StringRedisSerializer stringRedisSerializer = new StringRedisSerializer();
// 使用Jackson2JsonRedisSerialize 替换默认序列化(默认采用的是JDK序列化)
Jackson2JsonRedisSerializer<Object> jackson2JsonRedisSerializer = new Jackson2JsonRedisSerializer<>(Object.class);
ObjectMapper om = new ObjectMapper();
om.setVisibility(PropertyAccessor.ALL, JsonAutoDetect.Visibility.ANY);
om.enableDefaultTyping(ObjectMapper.DefaultTyping.NON_FINAL);
jackson2JsonRedisSerializer.setObjectMapper(om);
template.setValueSerializer(jackson2JsonRedisSerializer);
// key 采用 String的序列化
template.setKeySerializer(stringRedisSerializer);
// hash 的key采用String的序列化
template.setHashKeySerializer(stringRedisSerializer);
// hash 的 value 采用 String 的序列化
template.setHashValueSerializer(jackson2JsonRedisSerializer);
template.afterPropertiesSet();
return template;
}
@Bean
public RedisMessageListenerContainer container(RedisConnectionFactory factory, RedisMessageListener listener) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(factory);
// 订阅频道msgRedisTopic 这个container 可以添加多个 messageListener
container.addMessageListener(listener, new ChannelTopic(MessageHandler.CHANNEL_NAME));
return container;
}
}
RedisMessageListener:
@Component
public class RedisMessageListener implements MessageListener {
@Resource
private RedisTemplate<String, Object> redisTemplate;
@Override
public void onMessage(Message message, byte[] bytes) {
// 获取消息
byte[] messageBody = message.getBody();
MessageDto messageDto = Convert.convert(new TypeReference<MessageDto>() {
}, redisTemplate.getValueSerializer().deserialize(messageBody));
Map<String, WebSocketSession> onlineSessionMap = MessageHandler.CLIENTS;
if (onlineSessionMap.containsKey(messageDto.getToUid())) {
try {
onlineSessionMap.get(messageDto.getToUid()).sendMessage(new TextMessage("收到" messageDto.getSendUid() "的消息:" messageDto.getMsg()));
} catch (IOException e) {
e.printStackTrace();
}
}
}
}
WebSocketConfig:
@Configuration
public class WebSocketConfig implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(myHandler(), "/demo")
.addInterceptors(new MyInterceptor())
.setAllowedOrigins("*");
}
@Bean
public WebSocketHandler myHandler() {
return new MessageHandler();
}
}
MyInterceptor:
@Slf4j
@Component
public class MyInterceptor implements HandshakeInterceptor {
/**
* 握手前
*/
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception {
log.info("握手开始");
// 获得请求参数
Map<String, String> paramMap = HttpUtil.decodeParamMap(request.getURI().getQuery(), Charset.defaultCharset());
String uid = paramMap.get("uid");
if (CharSequenceUtil.isNotBlank(uid)) {
// 放入属性域
attributes.put("uid", uid);
log.info("用户{}握手成功!", uid);
return true;
}
return false;
}
/**
* 握手后
*/
@Override
public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception exception) {
log.info("握手完成");
}
}
MessageHandler:
@Component
public class MessageHandler extends TextWebSocketHandler {
@Resource
private RedisTemplate redisTemplate;
/**
* redis 订阅通道名
*/
public static final String CHANNEL_NAME = "msgRedisTopic";
/**
* 当前节点在线session
*/
public static Map<String, WebSocketSession> CLIENTS = new ConcurrentHashMap<>();
@Override
public void afterConnectionEstablished(WebSocketSession session) {
String uid = String.valueOf(session.getAttributes().get("uid"));
CLIENTS.put(uid, session);
log.info("uri :" session.getUri());
log.info("连接建立:uid{} ", uid);
log.info("当前连接服务器客户端数: {}", CLIENTS.size());
log.info("===================================");
}
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
String uid = String.valueOf(session.getAttributes().get("uid"));
CLIENTS.remove(uid);
log.info("断开连接: uid{}", uid);
log.info("当前连接服务器客户端数: {}", CLIENTS.size());
}
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws IOException {
String payload = message.getPayload();
log.info("服务端收到消息:{}", payload);
MessageDto messageDto = JSONUtil.toBean(payload, MessageDto.class);
String toUid = messageDto.getToUid();
if (CLIENTS.containsKey(toUid)) {
try {
log.info("当前ws服务器内包含客户端uid {},直接发送消息", toUid);
CLIENTS.get(toUid).sendMessage(new TextMessage("收到" messageDto.getSendUid() "的信息:" payload));
} catch (Exception e) {
e.printStackTrace();
}
} else {
log.warn("当前ws服务器内未找到客户端uid {},推送到redis", toUid);
// 发布消息
redisTemplate.convertAndSend(CHANNEL_NAME, messageDto);
}
}
}
测试
1.注意idea要开启允许一个服务启动多个实例
2.启动多个实例时记得改端口号
3.开多个ws测试工具
这篇好文章是转载于:学新通技术网
- 版权申明: 本站部分内容来自互联网,仅供学习及演示用,请勿用于商业和其他非法用途。如果侵犯了您的权益请与我们联系,请提供相关证据及您的身份证明,我们将在收到邮件后48小时内删除。
- 本站站名: 学新通技术网
- 本文地址: /boutique/detail/tanhffgfkk
系列文章
更多
同类精品
更多
-
微信小程序没声音怎么办
PHP中文网 06-15 -
怎样阻止微信小程序自动打开
PHP中文网 06-13 -
excel图片置于文字下方的方法
PHP中文网 06-27 -
微信人名旁边有个图标有什么用
PHP中文网 03-11 -
微信提示登录环境异常是什么意思原因
PHP中文网 04-09 -
微信获取用户openid失败怎么办
PHP中文网 03-26 -
photoshop怎么把印章抠出并放在另一张图上
PHP中文网 06-15 -
Excel筛选和排序是灰色的怎么办
PHP中文网 06-22 -
EhViewer(E绅士)最新版_ehviewer白色版彩色版_Ehviewer显示网络错误怎么办?e站进不去了怎么办
Evanpatchouli 09-19 -
photoshop蒙版画笔没反应怎么办
PHP中文网 06-24