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

Java如何从内部Runnable销毁Binance WebSocket连接类AccountStream实例

解决方案

你无法直接销毁实例停止运行的核心原因是:定时线程池、WebSocket客户端、Rest客户端都被定义为方法内部局部变量,类的其他位置无法获取引用执行关闭操作。你只需要将这些资源提升为类成员变量,统一封装销毁逻辑即可,修改后的完整代码如下:

public class AccountStream extends Driver {

    private Integer agentId;
    private String API_KEY;
    private String SECRET;
    private String listenKey;
    private Order newOrder;
    private String LOG_FILE_PATH;

    // 新增资源类成员变量,用于后续销毁操作
    private BinanceApiRestClient restClient;
    private BinanceApiWebSocketClient webSocketClient;
    private ScheduledExecutorService keepAlivePool;
    private ScheduledExecutorService connectionCheckPool;
    private volatile boolean isShutdown = false; // 状态位避免重复执行销毁逻辑

    public AccountStream(Integer agentId) {
        this.agentId = agentId;

        // Load binance config
        HashMap<String, String> binanceConfig = MainDriver.getBinanceConfig(agentId);
        API_KEY = binanceConfig.get("api_key");
        SECRET = binanceConfig.get("secret");

        startAccountEventStreaming();
        setConnectionCheckScheduler();
    }

    private void startAccountEventStreaming() {
        BinanceApiClientFactory factory = BinanceApiClientFactory.newInstance(API_KEY, SECRET);
        // 赋值给成员变量
        restClient = factory.newRestClient();

        // First, we obtain a listenKey which is required to interact with the user data stream
        listenKey = restClient.startUserDataStream();

        // Then, we open a new web socket client, and provide a callback that is called on every update
        // 赋值给成员变量
        webSocketClient = factory.newWebSocketClient();
        
        // Listen for changes in the account
        webSocketClient.onUserDataUpdateEvent(listenKey, response -> {
            System.out.println(response);
        });

        // Ping the datastream every 30 minutes to prevent a timeout
        // 赋值给成员变量
        keepAlivePool = Executors.newScheduledThreadPool(1);
        Runnable pingUserDataStream = () -> {
            if (!isShutdown) {
                restClient.keepAliveUserDataStream(listenKey);
            }
        };
        keepAlivePool.scheduleWithFixedDelay(pingUserDataStream, 0, 30, TimeUnit.MINUTES);
    }

    private void setConnectionCheckScheduler() {
        // 赋值给成员变量
        connectionCheckPool = Executors.newScheduledThreadPool(1);
        Runnable checkConnectionTask = () -> {
            if (!MainDriver.connected && !isShutdown) {
                // 调用销毁方法
                shutdown();
            }
        };
        connectionCheckPool.scheduleWithFixedDelay(checkConnectionTask, 0, 1, TimeUnit.SECONDS);
    }

    // 统一销毁逻辑
    public void shutdown() {
        if (isShutdown) {
            return;
        }
        isShutdown = true;

        // 1. 关闭所有定时线程池,终止待执行任务
        if (keepAlivePool != null && !keepAlivePool.isShutdown()) {
            keepAlivePool.shutdownNow();
        }
        if (connectionCheckPool != null && !connectionCheckPool.isShutdown()) {
            connectionCheckPool.shutdownNow();
        }

        // 2. 主动关闭Binance用户数据流
        if (restClient != null && listenKey != null) {
            try {
                restClient.closeUserDataStream(listenKey);
            } catch (Exception ignore) {
                // 断网状态下调用失败可直接忽略
            }
        }

        // 3. 关闭WebSocket连接
        if (webSocketClient != null) {
            webSocketClient.close();
        }

        // 4. 所有引用置空,方便GC回收当前实例
        restClient = null;
        webSocketClient = null;
        keepAlivePool = null;
        connectionCheckPool = null;
        API_KEY = null;
        SECRET = null;
        listenKey = null;
        agentId = null;
    }
}

修改后检测到断网时会自动执行shutdown()方法清理所有运行资源,之后你只需要将旧的AccountStream实例引用置空,重新创建新的实例即可重启数据流。


内容的提问来源于stack exchange,提问作者A. Vreeswijk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 19:24:01