1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89
| import com.banksteel.openerp.commons.framework.exception.EntityNotFoundException; import com.banksteel.openerp.commons.queue.MessageEntity; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.handler.TextWebSocketHandler;
import javax.websocket.server.ServerEndpoint; import java.util.Hashtable; import java.util.Map; import java.util.Timer;
@ServerEndpoint(value = "/api/websocket") public class WebsocketEndPoint extends TextWebSocketHandler {
private Timer timer;
private Logger logger = (Logger) LoggerFactory.getLogger(this.getClass()); private static Map<String, String> sessionIds = new Hashtable<String, String>(); private static Map<String, WebSocketSession> onlineSessions = new Hashtable<String, WebSocketSession>();
@Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { if (!session.isOpen()) { logger.info("获取session失败======="); timer.cancel(); return; } logger.info("获取session成功======="); logger.info("主动触发用户信息:{}", session.getUri().getQuery()); super.handleTextMessage(session, message); TextMessage returnMessage = new TextMessage(message.getPayload()); synchronized (WebsocketEndPoint.class) { if ("ping".equals(message.getPayload())) { TextMessage pong = new TextMessage("pong"); logger.info("{},{},{}", "心跳重连", "消息推送", String.valueOf(pong)); session.sendMessage(pong); } else { logger.info("{},{},{},{}", "主动触发", "消息推送", "首页消息完成", returnMessage); session.sendMessage(returnMessage); } } }
@Override public void afterConnectionEstablished(WebSocketSession session) { logger.info("===========用户正在连接============"); String userId = session.getUri().getQuery().split("=")[1]; logger.info("===========用户正在连接Id={}============", userId); if (userId != null && !userId.equals("")) { onlineSessions.put(userId, session); sessionIds.put(session.getId(), userId); } else { throw new EntityNotFoundException("未传入当前登录人id"); }
}
public void handleAllMessage(MessageEntity mess) throws Exception { for (String userId : mess.getUserIds()) { WebSocketSession session = onlineSessions.get(userId); if (session != null) { logger.info("MQ触发用户信息:{}", session.getUri().getQuery()); TextMessage textMessage = new TextMessage(mess.getMessage()); logger.info("{},{},{},{},{}", "MQ触发handleAllMessage", "消息推送", "弹窗消息推送", textMessage, userId); handleMessage(session, textMessage); } } }
@Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) { String userId = sessionIds.get(session.getId()); onlineSessions.remove(userId); sessionIds.remove(session.getId()); logger.info("用户websocket断开id={}", userId); } }
|