1. 依赖
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>
2.WebSocket
WebSocketMessage
public interface WebSocketMessage<T> {
void sendMsg(String id, T message);
}
AbstractWebSocket<T>
@Slf4j
public abstract class AbstractWebSocket<T> implements WebSocketMessage<T> {
protected Session session;
public static final CopyOnWriteArrayList<AbstractWebSocket<?>> allSessions = new CopyOnWriteArrayList<>();
public static final ConcurrentMap<String, List<AbstractWebSocket<?>>> sessionPool = new ConcurrentHashMap<>();
@OnOpen
public void onOpen(@PathParam("projectId") String projectId, Session session, EndpointConfig config){
log.info("create new session sessionId:{}", session.getId());
this.session = session;
allSessions.add(this);
synchronized (sessionPool){
List<AbstractWebSocket<?>> sessions = sessionPool.compute(projectId, (oldProjectId, processWebSockets) -> {
if (ObjectUtils.isEmpty(processWebSockets)){
processWebSockets = new ArrayList<>();
processWebSockets.add(this);
}else{
processWebSockets.add(this);
}
return processWebSockets;
});
sessionPool.put(projectId, sessions);
}
};
@OnClose
public void onClose(@PathParam("projectId") String projectId, Session session){
log.info("close session");
allSessions.remove(this);
sessionPool.computeIfPresent(projectId, (oldProjectId, processWebSockets) -> {
if(!ObjectUtils.isEmpty(processWebSockets)){
processWebSockets.remove(this);
}
if(processWebSockets.isEmpty()){
sessionPool.remove(projectId);
return null;
}else{
return processWebSockets;
}
});
}
@OnMessage
public void onMessage(@PathParam("projectId") String projectId,String message){
log.debug("message:{}, projectId:{}", message, projectId);
}
@OnError
public void onError(Session session, Throwable error){
log.error(error.getMessage(), error);
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
AbstractWebSocket<T> that = (AbstractWebSocket<T>) o;
return session.getId().equals(that.session.getId());
}
@Override
public int hashCode() {
return session.hashCode();
}
public Session getSession() {
return session;
}
public void setSession(Session session) {
this.session = session;
}
ProcessWebSocket
@Component
@Slf4j
@ServerEndpoint(WsVersion.workWs + "/process/{projectId}")
public class ProcessWebSocket extends AbstractWebSocket<Boolean> {
@Autowired
private ProjectService projectService;
@Override
public void sendMsg(String id, Boolean refreshComplete) {
ProcessWsResponse processWsResponse = new ProcessWsResponse();
List<VirtualSelectProcessQuery> virtualSelectProcessQueries = projectService.process(id, null);
List<VirtualSelectProcessQuery> errorArray = projectService.process(id, false);
processWsResponse.setCalculatingCount(virtualSelectProcessQueries.size());
processWsResponse.setCalculatingArray(virtualSelectProcessQueries);
processWsResponse.setErrorArray(errorArray);
processWsResponse.setErrorCount(errorArray.size());
processWsResponse.setRefreshComplete(refreshComplete);
processWsResponse.setCompleteCount(projectService.countByProjectId(id));
List<AbstractWebSocket<?>> webSockets = sessionPool.get(id);
if(!ObjectUtils.isEmpty(webSockets)){
for (AbstractWebSocket<?> abstractWebSocket : webSockets) {
try{
abstractWebSocket.getSession().getAsyncRemote().sendText(JSON.toJSONString(processWsResponse));
}catch (Exception e){
log.error(e.getMessage(), e);
}
}
}
}
@OnMessage
public void onMessage(@PathParam("projectId") String projectId, String message){
super.onMessage(projectId, message);
}
}
WebSocketConfig
@Configuration
@EnableWebSocket
public class WebSocketConfig {
@Bean
public ServerEndpointExporter serverEndpointExporter(){
return new ServerEndpointExporter();
}
}
这里需要注意每次连接都会创建一个AbstractWebSocket对应子对象的实例,因此使用了静态变量缓存了对应的session对象
K8s发布Websocket
apiVersion: extensions/v1beta1
kind: Ingress
metadata:
name: work-websocket
namespace: ${NAMESPACE}
annotations:
nginx.org/websocket-services: "work-api"
kubernetes.io/ingress.class: "nginx"
nginx.ingress.kubernetes.io/configuration-snippet: |
proxy_set_header Upgrade "websocket";
proxy_set_header Connection "Upgrade";
nginx.ingress.kubernetes.io/proxy-read-timeout: 3600s
nginx.ingress.kubernetes.io/proxy-write-timeout: 3600s
spec:
rules:
- host: xxx
http:
paths:
- path: /ws/v1/work
backend:
serviceName: work-api
servicePort: 8082