Java 21 WebSocket无法接收CoinEx服务器消息求助
我用Java 21编写了连接CoinEx期货WebSocket服务器(wss://socket.coinex.com/v2/futures)的代码,发送订阅消息后,onText方法始终无法接收到服务器返回的数据。但使用websocat发送相同的订阅消息可以正常接收响应。已知消息需要在接收后解压,但当前核心问题是onText方法根本没有收到任何数据,请求协助排查原因。
我的Java代码
public class MainClass { public static void main(String[] args) throws Exception { CountDownLatch latch = new CountDownLatch(1); try (HttpClient client = HttpClient.newHttpClient()) { client.newWebSocketBuilder() .buildAsync(URI.create("wss://socket.coinex.com/v2/futures"), new WebSocketClient(latch)) .join(); latch.wait(); } } private record WebSocketClient(CountDownLatch latch) implements WebSocket.Listener { @Override public void onOpen(WebSocket webSocket) { System.out.println("Connected to server"); String message = """ { "method": "state.subscribe", "params": {"market_list": ["BTCUSDT"]}, "id": 1 } """; webSocket.sendText(message, true); WebSocket.Listener.super.onOpen(webSocket); } @Override public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) { System.out.println("Receive: " + data.toString()); latch.countDown(); return WebSocket.Listener.super.onText(webSocket, data, last); } @Override public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) { System.out.println("Socket Closed: " + statusCode); latch.countDown(); return WebSocket.Listener.super.onClose(webSocket, statusCode, reason); } @Override public void onError(WebSocket webSocket, Throwable error) { System.out.println("Error: " + error.getMessage()); latch.countDown(); WebSocket.Listener.super.onError(webSocket, error); } } }
websocat测试命令(可正常接收响应)
websocat --uncompress-gzip --binary wss://socket.coinex.com/v2/futures # 连接后输入以下消息即可收到响应: {"method": "state.subscribe","params": {"market_list": ["BTCUSDT"]},"id": 1}
相关接口说明(翻译自CoinEx官方文档)
CoinEx期货WebSocket行情状态订阅接口,用于订阅指定交易对的实时行情状态数据。订阅时需发送符合格式的JSON消息,包含三个字段:
method:固定为state.subscribeparams:包含market_list数组,指定要订阅的交易对(如BTCUSDT)id:自定义的请求ID,用于标识请求
服务器返回的数据采用gzip压缩格式,需解压后才能解析。
排查方向及解决方案
1. 发送的JSON消息包含多余空白字符
你通过多行字符串定义的订阅消息包含大量缩进和换行,虽然JSON语法允许空白字符,但部分服务器对格式敏感,可能导致解析失败。websocat发送的是紧凑格式的JSON,建议将消息改为紧凑格式或去除多余空白:
String message = "{\"method\": \"state.subscribe\",\"params\": {\"market_list\": [\"BTCUSDT\"]},\"id\": 1}";
2. 未处理二进制消息
CoinEx服务器返回的是gzip压缩的二进制数据,而非文本帧,因此你的onText方法无法捕获到数据。需要重写onBinary方法来接收二进制帧并解压:
@Override public CompletionStage<?> onBinary(WebSocket webSocket, ByteBuffer data, boolean last) { // 将ByteBuffer转换为InputStream并解压 try (InputStream inputStream = new ByteArrayInputStream(data.array()); GZIPInputStream gzipIn = new GZIPInputStream(inputStream)) { String decompressedData = new String(gzipIn.readAllBytes(), StandardCharsets.UTF_8); System.out.println("Received decompressed data: " + decompressedData); latch.countDown(); } catch (IOException e) { System.err.println("Decompression failed: " + e.getMessage()); } return WebSocket.Listener.super.onBinary(webSocket, data, last); }
注意:需要导入java.util.zip.GZIPInputStream和java.io.ByteArrayInputStream。
3. 未处理消息发送的异常
webSocket.sendText()返回的CompletionStage<Void>可能会在发送失败时抛出异常,但你的代码未监听该阶段,导致发送失败的情况无法被察觉。建议添加监听来确认消息是否发送成功:
webSocket.sendText(message, true) .whenComplete((unused, throwable) -> { if (throwable != null) { System.err.println("Failed to send message: " + throwable.getMessage()); latch.countDown(); } else { System.out.println("Subscription message sent successfully"); } });
4. 检查WebSocket连接状态
可以在onOpen中添加日志确认连接是否成功建立,同时确保latch.wait()不会提前被触发(比如onError或onClose意外调用)。
内容的提问来源于stack exchange,提问作者mah454

