You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java WebSocket重连同步失效问题求助

并发重连WebSocket会话的同步问题

我尝试阻止重连方法的并发访问,试过synchronized(this)、方法级同步、同步所有相关方法,结果都没用。

我有个WebSocket服务器连接需要定期重连,在@ClientEndpoint类的@OnClose方法中会立即触发会话重连。原本期望@OnClose事件调用的方法会等待reconnectSession里的同步块执行完毕,但实际运行日志显示重连启动了两次,不是预期的一次。

实际运行日志

init ws connections ended 
RECCONECTING TO BINANCE WS  STARTED 
sessionWrapper instance com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa from thread Thread-21
RECCONECTING SESSION com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa started 
  session with id 0 closed with code NORMAL_CLOSURE with reason phrase 
sessionWrapper instance com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa from thread Thread-21
RECCONECTING SESSION com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa started 
RECCONECTING SESSION com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa ended 
RECCONECTING SESSION com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa ended 
RECCONECTING TO BINANCE WS  ENDED 

期望日志

RECCONECTING SESSION com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa started
RECCONECTING SESSION com.mycompany.binancereconnect.WsSessionWrapper@7c4b13aa ended 

复现代码示例

主类 TestSynchronization

package com.mycompany.binancereconnect;

import jakarta.websocket.ContainerProvider;
import jakarta.websocket.WebSocketContainer;
import java.net.URI;
import java.util.ArrayList;
import java.util.List;

/**
 *
 * @author utza
 */
public class TestSynchronization {

    public List<WsSessionWrapper> wsSessions = new ArrayList<>();

    public static String BASE_SPOT_SOCKET_URL = "wss://stream.binance.com:443/ws";

    public void init() {
        try {
            for (WsSessionWrapper sessionWrapper : wsSessions) {
                WebSocketContainer container = ContainerProvider.getWebSocketContainer();
                if (sessionWrapper.getType().equals("SPOT")) {
                    sessionWrapper.setSession(container.connectToServer(new BinanceSpotMessageReceiver(this),
                            new URI(BASE_SPOT_SOCKET_URL)));
                }

            }
            System.out.println(" init ws connections ended ");
        } catch (Exception ex) {
            ex.printStackTrace();

        }
    }

    public void reconnectSession(WsSessionWrapper sessionWrapper) {
        System.out.println("sessionWrapper instance " + sessionWrapper + " from thread " + Thread.currentThread().getName());
        synchronized (sessionWrapper) {
            if (sessionWrapper.allowedToReconnect()) {
                try {
                    System.out.println("RECCONECTING SESSION " + sessionWrapper + " started ");
                    try {
                        if (sessionWrapper.getSession() != null) {
                            sessionWrapper.getSession().close();
                        }
                    } catch (Exception ex) {
                        ex.printStackTrace();
                    }
                    sessionWrapper.setSession(null);
                    WebSocketContainer container = ContainerProvider.getWebSocketContainer();
                    if (sessionWrapper.getType().equals("SPOT")) {
                        sessionWrapper.setSession(container.connectToServer(new BinanceSpotMessageReceiver(this),
                                new URI(BASE_SPOT_SOCKET_URL)));
                    }
                    System.out.println("RECCONECTING SESSION " + sessionWrapper + " ended ");
                } catch (Exception ex) {
                    ex.printStackTrace();

                }
            }
        }
    }

    public void fullReconnect() {
//        aici;to also take into account no transactions on btc
        System.out.println("RECCONECTING TO BINANCE WS  STARTED ");
        for (WsSessionWrapper sessionWrapper : wsSessions) {
            reconnectSession(sessionWrapper);
        }

        System.out.println("RECCONECTING TO BINANCE WS  ENDED ");
    }

    public void reconnectSessionBySessionId(String sessionId) {
        for (WsSessionWrapper sessionWrapper : wsSessions) {
            if (sessionWrapper.getSession() != null && sessionId.equals(sessionWrapper.getSession().getId())) {
                reconnectSession(sessionWrapper);
            }
        }
    }

    public static void main(String[] args) throws InterruptedException {
        TestSynchronization testSynchronization = new TestSynchronization();
        WsSessionWrapper s1 = new WsSessionWrapper("SPOT");
        s1.setId("1");

        testSynchronization.wsSessions.add(s1);

        testSynchronization.init();
        Thread t1 = new Thread(() -> {
            testSynchronization.fullReconnect();
        });
        t1.start();
        Thread.sleep(100000);
    }

}

消息接收类 BinanceSpotMessageReceiver(@OnClose事件触发重连)

package com.mycompany.binancereconnect;


import jakarta.websocket.ClientEndpoint;
import jakarta.websocket.CloseReason;
import jakarta.websocket.OnClose;
import jakarta.websocket.OnMessage;
import jakarta.websocket.Session;
import java.io.IOException;
import java.util.Date;


/**
 *
 * @author utza
 */

@ClientEndpoint()
public class BinanceSpotMessageReceiver {

    public BinanceSpotMessageReceiver(TestSynchronization socketService) {
        this.socketService = socketService;
    }

    private TestSynchronization socketService;

    private Date lastTimeTradeEventReceived = new Date();

    @OnClose
    public void onClose(Session session, CloseReason closeReason) {
        if (closeReason != null) {
            System.out.println("  session with id " + session.getId() + " closed with code " + closeReason.getCloseCode() + " with reason phrase " + closeReason.getReasonPhrase());
        } else {
            System.out.println("  session with id " + session.getId() + " closed with close reason null  ");
        }
        socketService.reconnectSessionBySessionId(session.getId());
    }

    @OnMessage
    public void onMessage(Session session, String message) throws IOException {
        System.out.println("received message from binance  " + message);
    }
}

会话包装类 WsSessionWrapper

package com.mycompany.binancereconnect;

import jakarta.websocket.Session;
import java.util.Calendar;
import java.util.Date;

/**
 * 
 * @author utza
 */
public class WsSessionWrapper {

    private String id;
    private Session session;
    private Date lastReconnectionTime;
    private String type;

    public WsSessionWrapper(String type) {
        this.type = type;
    }

    public String getType() {
        return type;
    }

    public void setType(String type) {
        this.type = type;
    }

    public Session getSession() {
        return session;
    }

    public void setSession(Session session) {
        this.session = session;
    }

    public String getId() {
        return id;
    }

    public void setId(String id) {
        this.id = id;
    }

    public Date getLastReconnectionTime() {
        return lastReconnectionTime;
    }

    public void setLastReconnectionTime(Date lastReconnectionTime) {
        this.lastReconnectionTime = lastReconnectionTime;
    }

    //is allowed to reconnect if 1 minute passed from last reconnection
    public Boolean allowedToReconnect() {
        if (lastReconnectionTime == null || addMinutes(lastReconnectionTime, 1).compareTo(new Date()) < 0) {
            return true;
        }
        return false;
    }

    public static Date addMinutes(Date date, int minutes) {
        Calendar cal = Calendar.getInstance();
        cal.setTime(date);
        cal.add(Calendar.MINUTE, minutes);
        return cal.getTime();
    }
}

Maven依赖

<dependencies>
        <dependency>
            <groupId>jakarta.websocket</groupId>
            <artifactId>jakarta.websocket-client-api</artifactId>
            <version>2.1.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.tomcat.embed</groupId>
            <artifactId>tomcat-embed-websocket</artifactId>
            <version>10.1.7</version>
        </dependency>
        <dependency>
            <groupId>com.google.code.gson</groupId>
            <artifactId>gson</artifactId>
            <version>2.10.1</version>
        </dependency>
        <!-- https://mvnrepository.com/artifact/org.java-websocket/Java-WebSocket -->
        <dependency>
            <groupId>org.java-websocket</groupId>
            <artifactId>Java-WebSocket</artifactId>
            <version>1.5.1</version>
        </dependency>
    </dependencies>

恳请各位提供帮助,谢谢!

内容的提问来源于stack exchange,提问作者Uta Alexandru

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 06:42:00