赞
踩
springboot项目集成websocket服务,并且使用redis发布订阅方案解决websocket集群模式下session共享问题
<?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>
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
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测试工具
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。