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
相关产品推荐
相关产品推荐

