WebSocket的使用

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
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容