feat: 提交项目代码

This commit is contained in:
super_liu
2021-05-31 14:56:05 +08:00
parent d50bb875cb
commit 3ed9ca235a
234 changed files with 27792 additions and 0 deletions
@@ -0,0 +1,28 @@
package com.adc.da.websocket.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import javax.servlet.http.HttpSession;
import javax.websocket.HandshakeResponse;
import javax.websocket.server.HandshakeRequest;
import javax.websocket.server.ServerEndpointConfig;
import javax.websocket.server.ServerEndpointConfig.Configurator;
import java.util.Enumeration;
public class GetHttpSessionConfigurator extends Configurator {
private static final Logger logger = LoggerFactory.getLogger(GetHttpSessionConfigurator.class);
@Override
public void modifyHandshake(ServerEndpointConfig sec, HandshakeRequest request, HandshakeResponse response) {
HttpSession httpSession=(HttpSession) request.getHttpSession();
Enumeration<String> attributeNames = httpSession.getAttributeNames();
logger.info("---------------------------GetHttpSessionConfigurator--------------------------------");
while (attributeNames.hasMoreElements()){
String s = attributeNames.nextElement();
logger.info("建立session时的属性为:"+s+" 参数为:"+httpSession.getAttribute(s));
}
logger.info("---------------------------GetHttpSessionConfigurator--------------------------------");
sec.getUserProperties().put(HttpSession.class.getName(),httpSession);
}
}
@@ -0,0 +1,25 @@
package com.adc.da.websocket.config;
import com.adc.da.websocket.handler.SpringWebSocketHandler;
import com.adc.da.websocket.interceptor.SpringWebSocketHandlerInterceptor;
import org.springframework.web.servlet.config.annotation.WebMvcConfigurerAdapter;
import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry;
import org.springframework.web.socket.handler.TextWebSocketHandler;
/*@Configuration
@EnableWebSocket*/
public class SpringWebSocketConfig extends WebMvcConfigurerAdapter implements WebSocketConfigurer {
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
registry.addHandler(webSocketHandler(),"/websocket/socketServer").addInterceptors(new SpringWebSocketHandlerInterceptor()).setAllowedOrigins("*");
registry.addHandler(webSocketHandler(), "/sockjs/socketServer").addInterceptors(new SpringWebSocketHandlerInterceptor()).setAllowedOrigins("*").withSockJS();
}
// @Bean
public TextWebSocketHandler webSocketHandler(){
return new SpringWebSocketHandler();
}
}
@@ -0,0 +1,15 @@
package com.adc.da.websocket.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.server.standard.ServerEndpointExporter;
@Configuration
public class WebSocketConfig {
@Bean
public ServerEndpointExporter serverEndpointExporter() {
return new ServerEndpointExporter();
}
}
@@ -0,0 +1,26 @@
package com.adc.da.websocket.constant;
public enum SocketMesTypeEnum{
PROCESS_MESSAGE("process", "流程消息");
private String value;
private String text;
private SocketMesTypeEnum(String value, String text) {
this.value = value;
this.text = text;
}
public String getValue() {
return value;
}
public String getText() {
return text;
}
}
@@ -0,0 +1,39 @@
package com.adc.da.websocket.controller;
import com.adc.da.http.ResponseMessage;
import com.adc.da.http.Result;
import com.adc.da.websocket.entity.WebSocketMessage;
import com.adc.da.websocket.service.WebSocketServer;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.ArrayList;
import java.util.List;
@RestController
@RequestMapping("/${restPath}/webSocketCol")
@Api(description = "webSocket测试")
public class webSocketController {
@Autowired
// public WebSocketService webSocketService;
public WebSocketServer webSocketServer;
@ApiOperation(value = "|webSocket|发送消息给特定用户")
@GetMapping("/sendMessage")
public ResponseMessage sendMessage(String userId, WebSocketMessage message){
// webSocketService.sendMessageToUser(userId,message);
webSocketServer.sendMessage(userId,message);
List<String> userList =new ArrayList<String>();
userList.add(userId);
webSocketServer.sendMessageOfUserList(userList,message);
return Result.success("001","发送成功");
}
}
@@ -0,0 +1,57 @@
package com.adc.da.websocket.entity;
import java.util.concurrent.atomic.AtomicInteger;
import javax.websocket.Session;
/***
*webcocketBean对象
* date: 2018/12/6 18:26
*/
public class WebSocketBean {
/**
* 连接session对象
*/
private Session session;
/***
* socket绑定的用户信息
*/
private String userId;
/**
* 连接错误次数
*/
private AtomicInteger erroerLinkCount = new AtomicInteger(0);
public int getErroerLinkCount() {
// 线程安全,以原子方式将当前值加1,注意:这里返回的是自增前的值
return erroerLinkCount.getAndIncrement();
}
public void cleanErrorNum()
{
// 清空计数
erroerLinkCount = new AtomicInteger(0);
}
public Session getSession() {
return session;
}
public void setSession(Session session) {
this.session = session;
}
public void setErroerLinkCount(AtomicInteger erroerLinkCount) {
this.erroerLinkCount = erroerLinkCount;
}
public String getUserId() {
return userId;
}
public void setUserId(String userId) {
this.userId = userId;
}
}
@@ -0,0 +1,36 @@
package com.adc.da.websocket.entity;
import java.io.Serializable;
public class WebSocketMessage implements Serializable{
private String message;
private String type;
private String businessId;
public String getMessage() {
return message;
}
public void setMessage(String message) {
this.message = message;
}
public String getBusinessId() {
return businessId;
}
public void setBusinessId(String businessId) {
this.businessId = businessId;
}
public String getType() {
return type;
}
public void setType(String type) {
this.type = type;
}
}
@@ -0,0 +1,121 @@
package com.adc.da.websocket.handler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.web.socket.*;
import org.springframework.web.socket.handler.TextWebSocketHandler;
import java.io.IOException;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
public class SpringWebSocketHandler extends TextWebSocketHandler {
private static final Map<String, WebSocketSession> users;//这个会出现性能问题,最好用Map来存储,key用userid
private static final Logger logger = LoggerFactory.getLogger(SpringWebSocketHandler.class);
static {
users = new HashMap<String, WebSocketSession>();
}
public SpringWebSocketHandler() {
super();
}
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
// TODO Auto-generated method stub
System.out.println("connect to the websocket success......当前数量:" + users.size());
//此处获取用户session中的userId
String userId = (String) session.getAttributes().get("WEBSOCKET_USERID");
users.put(userId, session);
//这块会实现自己业务,比如,当用户登录后,会把离线消息推送给用户
//TextMessage returnMessage = new TextMessage("你将收到的离线");
//session.sendMessage(returnMessage);
super.afterConnectionEstablished(session);
}
@Override
public void handleMessage(WebSocketSession session, WebSocketMessage<?> message) throws Exception {
String userId = (String) session.getAttributes().get("WEBSOCKET_USERID");
logger.info(userId+": 发送了:"+message.toString());
super.handleMessage(session, message);
}
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
super.handleTextMessage(session, message);
}
@Override
protected void handlePongMessage(WebSocketSession session, PongMessage message) throws Exception {
super.handlePongMessage(session, message);
}
@Override
public void handleTransportError(WebSocketSession session, Throwable exception) throws Exception {
super.handleTransportError(session, exception);
}
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception {
logger.debug("websocket connection closed......");
String userId = (String) session.getAttributes().get("WEBSOCKET_USERID");
logger.info("用户" + userId + "已退出!");
users.remove(userId);
logger.info("剩余在线用户" + users.size());
}
@Override
public boolean supportsPartialMessages() {
return super.supportsPartialMessages();
}
/**
* 给某个用户发送消息
*
* @param userId
* @param message
*/
public void sendMessageToUser(String userId, TextMessage message) {
logger.info("----------------------------------webSocket-----------------");
if (users.containsKey(userId)) {
WebSocketSession user = users.get(userId);
if (user.isOpen()) {
try {
user.sendMessage(message);
} catch (IOException e) {
logger.error(""+userId+"用户推送消息时发生异常:" + e.getMessage(), e);
}
}
}
}
/**
* 给某个用户发送消息
*
* @param userIdList
* @param message
*/
public void sendMessageToUsers(List<String> userIdList, TextMessage message) {
if (userIdList != null && !userIdList.isEmpty()) {
for (String userId : userIdList) {
if (users.containsKey(userId)) {
WebSocketSession user = users.get(userId);
if (user.isOpen()) {
try {
user.sendMessage(message);
} catch (IOException e) {
logger.error(""+userId+"用户推送消息时发生异常:" + e.getMessage(), e);
}
}
}
}
}
}
}
@@ -0,0 +1,36 @@
package com.adc.da.websocket.interceptor;
import org.springframework.http.server.ServerHttpRequest;
import org.springframework.http.server.ServerHttpResponse;
import org.springframework.http.server.ServletServerHttpRequest;
import org.springframework.web.socket.WebSocketHandler;
import org.springframework.web.socket.server.support.HttpSessionHandshakeInterceptor;
import javax.servlet.http.HttpSession;
import java.util.Map;
public class SpringWebSocketHandlerInterceptor extends HttpSessionHandshakeInterceptor {
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception {
// TODO Auto-generated method stub
if (request instanceof ServletServerHttpRequest) {
ServletServerHttpRequest servletRequest = (ServletServerHttpRequest) request;
HttpSession session = servletRequest.getServletRequest().getSession(false);
if (session != null) {
//使用userName区分WebSocketHandler,以便定向发送消息
String userId = (String) session.getAttribute("LOGIN_USER_ID");
if (userId==null) {
userId="default-system";
}
attributes.put("WEBSOCKET_USERID",userId);
}
}
return super.beforeHandshake(request, response, wsHandler, attributes);
}
@Override
public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception ex) {
super.afterHandshake(request, response, wsHandler, ex);
}
}
@@ -0,0 +1,59 @@
package com.adc.da.websocket.service;
import com.adc.da.websocket.entity.WebSocketMessage;
import javax.websocket.EndpointConfig;
import javax.websocket.Session;
import java.util.List;
public interface WebSocketServer {
/**
* 连接建立成功调用的方法
* @param session session 对象
*/
public void onOpen(Session session, EndpointConfig config);
/**
* 断开连接方法
*/
public void onClose(Session session);
/**
* 收到客户端消息后调用的方法
* @param session session 对象
* @param message 返回客户端的消息
*/
public void onMessage(Session session, String message);
/**
* 发生异常时触发的方法
* @param session session 对象
* @param throwable 抛出的异常
*/
public void onError(Session session,Throwable throwable);
/**
* 向单个客户端发送消息
* @param userId 用户ID对象
* @param message 发送给客户端的消息
*/
public void sendMessage(String userId, WebSocketMessage message);
/***
* 向多个用户发送消息
* @MethodName:sendMessageOfUserList
* @author: zhangyanduan
* @param:[userList, message]
* @return:void
* date: 2018/12/17 15:04
*/
public void sendMessageOfUserList(List<String> userList,WebSocketMessage message);
/**
* 向所有在线用户群发消息
* @param message 发送给客户端的消息
*/
public void batchSendMessage(String message);
}
@@ -0,0 +1,44 @@
package com.adc.da.websocket.service;
import com.adc.da.websocket.handler.SpringWebSocketHandler;
import org.springframework.web.socket.TextMessage;
import java.util.List;
//@Service("webSocketService")
public class WebSocketService {
// @Bean//这个注解会从Spring容器拿出Bean
public SpringWebSocketHandler infoHandler() {
return new SpringWebSocketHandler();
}
/***
* 向某用户推送消息
* @MethodName:sendMessageToUser
* @author: DuYunbao
* @param:[userId, messageText]
* @return:void
* date: 2018/11/30 14:49
*/
public void sendMessageToUser(String userId,String messageText){
// infoHandler().sendMessageToUser(userId,new TextMessage(messageText));
}
/***
* 同时向多个用户推送统一消息
* @MethodName:sendMessageToUserList
* @author: DuYunbao
* @param:[userList, messageText]
* @return:void
* date: 2018/11/30 14:50
*/
public void sendMessageToUserList(List<String> userList,String messageText){
infoHandler().sendMessageToUsers(userList,new TextMessage(messageText));
}
}
@@ -0,0 +1,140 @@
package com.adc.da.websocket.service;
import com.adc.da.websocket.config.GetHttpSessionConfigurator;
import com.adc.da.websocket.entity.WebSocketBean;
import com.adc.da.websocket.entity.WebSocketMessage;
import com.alibaba.fastjson.JSONObject;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import javax.servlet.http.HttpSession;
import javax.websocket.*;
import javax.websocket.server.PathParam;
import javax.websocket.server.ServerEndpoint;
import java.util.Enumeration;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@ServerEndpoint(value = "/websocket/socketServer/{userId}",configurator= GetHttpSessionConfigurator.class)
@Component("webSocketService")
public class WebsocketServerImpl implements WebSocketServer {
private static final Logger logger = LoggerFactory.getLogger(WebsocketServerImpl.class);
/**
* 错误最大重试次数
*/
private static final int MAX_ERROR_NUM = 10;
/**
* 用来存放每个客户端对应的webSocket对象。
*/
private static Map<String, WebSocketBean> webSocketInfo;
static
{
// concurrent包的线程安全map
webSocketInfo = new ConcurrentHashMap<String, WebSocketBean>();
}
@OnOpen
public void onOpen(@PathParam(value = "userId") String userId, Session session, EndpointConfig config) {
// 如果是session没有激活的情况,就是没有请求获取或session,这里可能会取出空,需要实际业务处理
HttpSession httpSession= (HttpSession) config.getUserProperties().get(HttpSession.class.getName());
if(httpSession != null)
{
logger.info("获取到httpsession" + httpSession.getId());
// String userId = (String) httpSession.getAttribute(RequestUtils.LOGIN_USER_ID);
Enumeration<String> attributeNames = httpSession.getAttributeNames();
while(attributeNames.hasMoreElements()){
String s = attributeNames.nextElement();
logger.info("建立session时的属性为:"+s+" 参数为:"+httpSession.getAttribute(s));
}
logger.info("--------websocket 打开链接 当前用户为:"+userId);
}else {
logger.error("未获取到httpsession");
}
// 连接成功当前对象放入websocket对象集合
WebSocketBean bean = new WebSocketBean();
bean.setSession(session);
webSocketInfo.put(userId,bean);
logger.info("客户端连接服务器session id :"+session.getId()+",当前连接数:" + webSocketInfo.size());
}
@Override
public void onOpen(Session session, EndpointConfig config) {
logger.info("-------------------------此处仅为实现其接口即可-----------------------------------");
}
@OnClose
@Override
public void onClose(Session session) {
// 客户端断开连接移除websocket对象
webSocketInfo.remove(session.getId());
logger.info("客户端断开连接,当前连接数:" + webSocketInfo.size());
}
@OnMessage
@Override
public void onMessage(Session session, String message) {
logger.info("客户端 session id: "+session.getId()+",消息:" + message);
// 此方法为客户端给服务器发送消息后进行的处理,可以根据业务自己处理,这里返回页面
// sendMessage(session, "服务端返回" + message);
}
@OnError
@Override
public void onError(Session session, Throwable throwable) {
logger.error("发生错误"+ throwable.getMessage(),throwable);
}
@Override
public void sendMessage(String userId, WebSocketMessage webSocketMessage) {
String message = JSONObject.toJSONString(webSocketMessage);
try {
WebSocketBean webSocketBean = webSocketInfo.get(userId);
if(webSocketBean != null){
webSocketBean.getSession().getBasicRemote().sendText(message);
webSocketBean.cleanErrorNum();
}else{
logger.info("用户:"+userId+" 当前不在线!");
}
} catch (Exception e) {
logger.error("发送消息失败"+ e.getMessage(),e);
try {
int errorNum = webSocketInfo.get(userId).getErroerLinkCount();
// 小于最大重试次数重发
if (errorNum <= MAX_ERROR_NUM) {
sendMessage(userId, webSocketMessage);
} else {
logger.error("发送消息失败超过最大次数");
// 清空错误计数
webSocketInfo.get(userId).cleanErrorNum();
}
}catch (Exception e1){
logger.error(e1.getMessage(),e1);
}
}
}
@Override
public void sendMessageOfUserList(List<String> userList, WebSocketMessage message) {
if(userList !=null && !userList.isEmpty()){
for(String userId:userList){
sendMessage(userId, message);
}
}
}
@Override
public void batchSendMessage(String message) {
// 不做操作
}
}