> For the complete documentation index, see [llms.txt](https://docs.infoway.io/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://docs.infoway.io/websocket-api/code-examples.md).

# Websocket代码示例

WebSocket代码示例，获取实时成交明细、盘口、k线。包括自动重连，心跳检测机制。

{% tabs %}
{% tab title="Java" %}

```java
import com.alibaba.fastjson2.JSONArray;
import com.alibaba.fastjson2.JSONObject;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import jakarta.websocket.ClientEndpoint;
import jakarta.websocket.CloseReason;
import jakarta.websocket.ContainerProvider;
import jakarta.websocket.OnClose;
import jakarta.websocket.OnError;
import jakarta.websocket.OnMessage;
import jakarta.websocket.OnOpen;
import jakarta.websocket.Session;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;

import java.io.IOException;
import java.net.URI;
import java.net.URLEncoder;
import java.nio.charset.StandardCharsets;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;

@ClientEndpoint
@Slf4j
@Component
public class WebsocketExample {

    private static final Set<Integer> ACK_CODES = Set.of(10001, 10004, 10007, 11010);

    private volatile Session session;
    private ScheduledExecutorService scheduler;
    private volatile ScheduledFuture<?> heartbeatFuture;
    private volatile boolean authFailed;

    private String wsUrl() {
        String key = System.getenv().getOrDefault("INFOWAY_API_KEY", "yourApikey");
        return "wss://data.infoway.io/ws?business=crypto&apikey="
                + URLEncoder.encode(key, StandardCharsets.UTF_8);
    }

    @PostConstruct
    public void connectAll() {
        // 订阅与心跳不要挤在同一个单线程上
        scheduler = Executors.newScheduledThreadPool(2);
        connect();
        scheduler.scheduleAtFixedRate(() -> {
            if (authFailed) {
                return;
            }
            Session current = session;
            if (current == null || !current.isOpen()) {
                connect();
            }
        }, 10, 10, TimeUnit.SECONDS);
    }

    private void connect() {
        try {
            cancelHeartbeat();
            if (session != null && session.isOpen()) {
                session.close(new CloseReason(CloseReason.CloseCodes.NORMAL_CLOSURE, "reconnect"));
            }
            session = ContainerProvider.getWebSocketContainer().connectToServer(this, URI.create(wsUrl()));
        } catch (Exception e) {
            session = null;
            String msg = String.valueOf(e.getMessage());
            if (msg.contains("401") || msg.contains("407") || msg.contains("507")) {
                authFailed = true;
                log.error("鉴权失败，停止重连: {}", msg);
                return;
            }
            log.error("WebSocket 连接失败: {}", msg);
        }
    }

    @OnOpen
    public void onOpen(Session session) {
        this.session = session;
        authFailed = false;
        startHeartbeat(session);
        scheduler.execute(() -> {
            try {
                sendTradeSubscribe(session);
                sendDepthSubscribe(session);
                sendKlineSubscribe(session);
            } catch (IOException e) {
                log.error("发送订阅失败: {}", e.getMessage());
            }
        });
    }

    private void sendTradeSubscribe(Session session) throws IOException {
        JSONObject body = new JSONObject();
        body.put("code", 10000);
        body.put("trace", UUID.randomUUID().toString());
        JSONObject data = new JSONObject();
        data.put("codes", "BTCUSDT");
        body.put("data", data);
        session.getBasicRemote().sendText(body.toJSONString());
    }

    private void sendDepthSubscribe(Session session) throws IOException {
        JSONObject body = new JSONObject();
        body.put("code", 10003);
        body.put("trace", UUID.randomUUID().toString());
        JSONObject data = new JSONObject();
        data.put("codes", "BTCUSDT");
        body.put("data", data);
        session.getBasicRemote().sendText(body.toJSONString());
    }

    private void sendKlineSubscribe(Session session) throws IOException {
        JSONObject body = new JSONObject();
        body.put("code", 10006);
        body.put("trace", UUID.randomUUID().toString());
        JSONObject item = new JSONObject();
        item.put("type", 1);
        item.put("codes", "BTCUSDT");
        JSONObject data = new JSONObject();
        data.put("arr", new JSONArray().fluentAdd(item));
        body.put("data", data);
        session.getBasicRemote().sendText(body.toJSONString());
    }

    private void startHeartbeat(Session active) {
        cancelHeartbeat();
        heartbeatFuture = scheduler.scheduleAtFixedRate(() -> {
            try {
                if (active.isOpen() && active == this.session) {
                    JSONObject ping = new JSONObject();
                    ping.put("code", 10010);
                    ping.put("trace", UUID.randomUUID().toString());
                    active.getBasicRemote().sendText(ping.toJSONString());
                }
            } catch (IOException e) {
                log.error("发送心跳失败: {}", e.getMessage());
            }
        }, 30, 30, TimeUnit.SECONDS);
    }

    private void cancelHeartbeat() {
        ScheduledFuture<?> future = heartbeatFuture;
        if (future != null) {
            future.cancel(false);
            heartbeatFuture = null;
        }
    }

    @OnMessage
    public void onMessage(String message, Session session) {
        JSONObject msg = JSONObject.parseObject(message);
        if (msg == null || !msg.containsKey("code")) {
            log.info("忽略非 JSON 帧: {}", message);
            return;
        }
        Integer code = msg.getInteger("code");
        if (code == null) {
            return;
        }
        Object data = msg.get("data");
        switch (code) {
            case 200 -> log.info("连接成功: {}", msg.getString("msg"));
            case 10002 -> log.info("成交推送: {}", data);
            case 10005 -> log.info("盘口推送: {}", data);
            case 10008 -> log.info("K线推送: {}", data);
            case 10011 -> log.debug("心跳确认: {}", msg.getString("trace"));
            default -> {
                if (ACK_CODES.contains(code)) {
                    log.info("订阅确认 code={} trace={} msg={}", code, msg.getString("trace"), msg.getString("msg"));
                } else if (code >= 500 && code < 10000) {
                    log.error("业务错误 code={} msg={}", code, msg.getString("msg"));
                    if (code == 507) {
                        authFailed = true;
                    }
                } else {
                    log.warn("未处理协议号 {}: {}", code, message);
                }
            }
        }
    }

    @OnClose
    public void onClose(Session session, CloseReason reason) {
        cancelHeartbeat();
        this.session = null;
    }

    @OnError
    public void onError(Session session, Throwable error) {
        log.error("WebSocket 错误", error);
        cancelHeartbeat();
        this.session = null;
    }

    @PreDestroy
    public void destroy() {
        cancelHeartbeat();
        if (session != null && session.isOpen()) {
            try {
                session.close(new CloseReason(CloseReason.CloseCodes.NORMAL_CLOSURE, "shutdown"));
            } catch (IOException ignored) {
            }
        }
        if (scheduler != null) {
            scheduler.shutdownNow();
        }
    }
}
```

{% endtab %}

{% tab title="Python" %}

```python
import asyncio
import json
import logging
import os
import uuid
from typing import Optional
from urllib.parse import quote

import websockets
from websockets.asyncio.client import ClientConnection
from websockets.exceptions import ConnectionClosed, InvalidStatus

logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
logger = logging.getLogger("infoway-ws")

REQ_TRADE, REQ_DEPTH, REQ_KLINE, REQ_HEARTBEAT = 10000, 10003, 10006, 10010
PUSH_TRADE, PUSH_DEPTH, PUSH_KLINE = 10002, 10005, 10008
ACK_CODES = {10001, 10004, 10007, 11010}
WELCOME, HEART_ACK = 200, 10011


class CryptoWebsocketClient:
    def __init__(self, api_key: str, business: str = "crypto"):
        self.ws_url = (
            f"wss://data.infoway.io/ws?business={quote(business)}"
            f"&apikey={quote(api_key)}"
        )
        self.ws: Optional[ClientConnection] = None
        self.running = True
        self.auth_failed = False
        self.heartbeat_task: Optional[asyncio.Task] = None

    async def _send(self, msg: dict) -> None:
        assert self.ws is not None
        await self.ws.send(json.dumps(msg, separators=(",", ":")))

    async def _subscribe(self) -> None:
        await self._send({"code": REQ_TRADE, "trace": str(uuid.uuid4()), "data": {"codes": "BTCUSDT"}})
        await self._send({"code": REQ_DEPTH, "trace": str(uuid.uuid4()), "data": {"codes": "BTCUSDT"}})
        await self._send({
            "code": REQ_KLINE,
            "trace": str(uuid.uuid4()),
            "data": {"arr": [{"type": 1, "codes": "BTCUSDT"}]},
        })

    def _start_heartbeat(self) -> None:
        self._cancel_heartbeat()

        async def loop():
            try:
                while True:
                    await asyncio.sleep(30)
                    if self.ws is None:
                        return
                    await self._send({"code": REQ_HEARTBEAT, "trace": str(uuid.uuid4())})
            except (ConnectionClosed, asyncio.CancelledError):
                return

        self.heartbeat_task = asyncio.create_task(loop())

    def _cancel_heartbeat(self) -> None:
        if self.heartbeat_task and not self.heartbeat_task.done():
            self.heartbeat_task.cancel()
        self.heartbeat_task = None

    def _on_message(self, raw) -> None:
        try:
            msg = json.loads(raw)
        except json.JSONDecodeError:
            logger.info("忽略非 JSON 帧: %s", raw)
            return
        code = msg.get("code")
        data = msg.get("data")
        if code == WELCOME:
            logger.info("连接成功: %s", msg.get("msg"))
        elif code == PUSH_TRADE:
            logger.info("成交推送: %s", data)
        elif code == PUSH_DEPTH:
            logger.info("盘口推送: %s", data)
        elif code == PUSH_KLINE:
            logger.info("K线推送: %s", data)
        elif code in ACK_CODES:
            logger.info("订阅确认 code=%s trace=%s msg=%s", code, msg.get("trace"), msg.get("msg"))
        elif code == HEART_ACK:
            logger.debug("心跳确认 trace=%s", msg.get("trace"))
        elif isinstance(code, int) and 500 <= code < 10000:
            logger.error("业务错误 code=%s msg=%s", code, msg.get("msg"))
            if code == 507:
                self.auth_failed = True
                self.running = False
        else:
            logger.warning("未处理协议号 %s: %s", code, raw)

    async def _connect_once(self) -> None:
        async with websockets.connect(self.ws_url) as ws:
            self.ws = ws
            await self._subscribe()
            self._start_heartbeat()
            try:
                async for raw in ws:
                    self._on_message(raw)
            finally:
                self._cancel_heartbeat()
                self.ws = None

    async def start(self) -> None:
        backoff = 5
        while self.running and not self.auth_failed:
            try:
                await self._connect_once()
                backoff = 5
            except InvalidStatus as e:
                logger.error("握手失败: %s", e)
                if getattr(e, "response", None) is not None and e.response.status_code in (401, 407):
                    self.auth_failed = True
                    break
            except ConnectionClosed as e:
                logger.warning("连接关闭: %s", e)
                backoff = 5
            except Exception as e:
                logger.error("连接异常: %s", e)
            if not self.running or self.auth_failed:
                break
            await asyncio.sleep(backoff)
            backoff = min(backoff * 2, 60)


async def main():
    api_key = os.environ.get("INFOWAY_API_KEY", "yourApikey")
    await CryptoWebsocketClient(api_key).start()


if __name__ == "__main__":
    asyncio.run(main())
```

{% endtab %}

{% tab title="Go" %}

```go
package main

import (
	"context"
	"encoding/json"
	"fmt"
	"log"
	"net/url"
	"os"
	"os/signal"
	"sync"
	"syscall"
	"time"

	"github.com/google/uuid"
	"github.com/gorilla/websocket"
)

type frame struct {
	Code  int             `json:"code"`
	Trace string          `json:"trace"`
	Msg   string          `json:"msg"`
	Data  json.RawMessage `json:"data"`
}

type Client struct {
	wsURL    string
	conn     *websocket.Conn
	writeMu  sync.Mutex
	done     chan struct{}
	authFail bool
}

func NewClient(apiKey, business string) *Client {
	q := url.Values{}
	q.Set("business", business)
	q.Set("apikey", apiKey)
	return &Client{
		wsURL: "wss://data.infoway.io/ws?" + q.Encode(),
		done:  make(chan struct{}),
	}
}

func (c *Client) writeJSON(v any) error {
	c.writeMu.Lock()
	defer c.writeMu.Unlock()
	if c.conn == nil {
		return fmt.Errorf("not connected")
	}
	return c.conn.WriteJSON(v)
}

func (c *Client) subscribe() {
	_ = c.writeJSON(map[string]any{"code": 10000, "trace": uuid.NewString(), "data": map[string]string{"codes": "BTCUSDT"}})
	_ = c.writeJSON(map[string]any{"code": 10003, "trace": uuid.NewString(), "data": map[string]string{"codes": "BTCUSDT"}})
	_ = c.writeJSON(map[string]any{
		"code":  10006,
		"trace": uuid.NewString(),
		"data":  map[string]any{"arr": []map[string]any{{"type": 1, "codes": "BTCUSDT"}}},
	})
}

func (c *Client) handle(raw []byte) {
	var msg frame
	if err := json.Unmarshal(raw, &msg); err != nil {
		log.Printf("忽略非 JSON 帧: %s", raw)
		return
	}
	switch msg.Code {
	case 200:
		log.Printf("连接成功: %s", msg.Msg)
	case 10002:
		log.Printf("成交推送: %s", msg.Data)
	case 10005:
		log.Printf("盘口推送: %s", msg.Data)
	case 10008:
		log.Printf("K线推送: %s", msg.Data)
	case 10001, 10004, 10007, 11010:
		log.Printf("订阅确认 code=%d trace=%s msg=%s", msg.Code, msg.Trace, msg.Msg)
	case 10011:
		log.Printf("心跳确认 trace=%s", msg.Trace)
	default:
		if msg.Code >= 500 && msg.Code < 10000 {
			log.Printf("业务错误 code=%d msg=%s", msg.Code, msg.Msg)
			if msg.Code == 507 {
				c.authFail = true
			}
			return
		}
		log.Printf("未处理协议号 %d: %s", msg.Code, raw)
	}
}

func (c *Client) runSession(ctx context.Context) error {
	conn, resp, err := websocket.DefaultDialer.Dial(c.wsURL, nil)
	if err != nil {
		if resp != nil && (resp.StatusCode == 401 || resp.StatusCode == 407) {
			c.authFail = true
		}
		return err
	}
	defer conn.Close()
	c.conn = conn
	c.subscribe()

	errCh := make(chan error, 1)
	go func() {
		for {
			_, raw, readErr := conn.ReadMessage()
			if readErr != nil {
				errCh <- readErr
				return
			}
			c.handle(raw)
		}
	}()

	ticker := time.NewTicker(30 * time.Second)
	defer ticker.Stop()
	for {
		select {
		case <-ctx.Done():
			return nil
		case <-c.done:
			return nil
		case err := <-errCh:
			return err
		case <-ticker.C:
			if err := c.writeJSON(map[string]any{"code": 10010, "trace": uuid.NewString()}); err != nil {
				return err
			}
		}
	}
}

func (c *Client) Start() {
	ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
	defer stop()
	backoff := 5 * time.Second
	for {
		if c.authFail {
			log.Print("鉴权失败，停止重连")
			return
		}
		err := c.runSession(ctx)
		if ctx.Err() != nil {
			return
		}
		if err != nil {
			log.Printf("连接结束: %v", err)
		}
		if c.authFail {
			return
		}
		select {
		case <-ctx.Done():
			return
		case <-time.After(backoff):
		}
		if backoff < 60*time.Second {
			backoff *= 2
		}
	}
}

func main() {
	apiKey := os.Getenv("INFOWAY_API_KEY")
	if apiKey == "" {
		apiKey = "yourApikey"
	}
	NewClient(apiKey, "crypto").Start()
}
```

{% endtab %}

{% tab title="Php" %}

```php
<?php

require __DIR__ . '/vendor/autoload.php';

use WebSocket\Client;
use WebSocket\ConnectionException;

class CryptoWebsocketClient
{
    private string $wsUrl;
    private ?Client $client = null;
    private bool $connected = false;
    private bool $running = true;
    private bool $authFailed = false;
    private int $lastHeartbeat = 0;

    public function __construct(string $apiKey, string $business = 'crypto')
    {
        $this->wsUrl = 'wss://data.infoway.io/ws?business=' . rawurlencode($business)
            . '&apikey=' . rawurlencode($apiKey);
    }

    private function send(array $msg): void
    {
        if (!$this->connected || !$this->client) {
            return;
        }
        $this->client->send(json_encode($msg, JSON_UNESCAPED_UNICODE));
    }

    private function subscribe(): void
    {
        $this->send(['code' => 10000, 'trace' => $this->trace(), 'data' => ['codes' => 'BTCUSDT']]);
        $this->send(['code' => 10003, 'trace' => $this->trace(), 'data' => ['codes' => 'BTCUSDT']]);
        $this->send([
            'code' => 10006,
            'trace' => $this->trace(),
            'data' => ['arr' => [['type' => 1, 'codes' => 'BTCUSDT']]],
        ]);
    }

    private function handle(string $raw): void
    {
        $msg = json_decode($raw, true);
        if (!is_array($msg) || !isset($msg['code'])) {
            echo date('c') . " 忽略非 JSON 帧: {$raw}\n";
            return;
        }
        $code = (int) $msg['code'];
        $data = $msg['data'] ?? null;
        switch ($code) {
            case 200:
                echo date('c') . " 连接成功: " . ($msg['msg'] ?? '') . "\n";
                break;
            case 10002:
                echo date('c') . " 成交推送: " . json_encode($data, JSON_UNESCAPED_UNICODE) . "\n";
                break;
            case 10005:
                echo date('c') . " 盘口推送: " . json_encode($data, JSON_UNESCAPED_UNICODE) . "\n";
                break;
            case 10008:
                echo date('c') . " K线推送: " . json_encode($data, JSON_UNESCAPED_UNICODE) . "\n";
                break;
            case 10001:
            case 10004:
            case 10007:
            case 11010:
                echo date('c') . " 订阅确认 code={$code} msg=" . ($msg['msg'] ?? '') . "\n";
                break;
            case 10011:
                break;
            default:
                if ($code >= 500 && $code < 10000) {
                    echo date('c') . " 业务错误 code={$code} msg=" . ($msg['msg'] ?? '') . "\n";
                    if ($code === 507) {
                        $this->authFailed = true;
                        $this->running = false;
                    }
                } else {
                    echo date('c') . " 未处理协议号 {$code}: {$raw}\n";
                }
        }
    }

    public function start(): void
    {
        $nextReconnect = 0;
        while ($this->running && !$this->authFailed) {
            if (!$this->connected) {
                if (time() < $nextReconnect) {
                    usleep(200000);
                    continue;
                }
                try {
                    $this->client = new Client($this->wsUrl, ['timeout' => 30]);
                    $this->connected = true;
                    $this->lastHeartbeat = time();
                    $this->subscribe();
                } catch (Exception $e) {
                    $this->connected = false;
                    if (str_contains($e->getMessage(), '401') || str_contains($e->getMessage(), '507')) {
                        $this->authFailed = true;
                        echo date('c') . " 鉴权失败，停止重连\n";
                        break;
                    }
                    echo date('c') . " 连接失败: {$e->getMessage()}\n";
                    $nextReconnect = time() + 10;
                    continue;
                }
            }

            try {
                $socket = $this->client->getSocket();
                stream_set_blocking($socket, false);
                $raw = $this->client->receive();
                if (is_string($raw) && $raw !== '') {
                    $this->handle($raw);
                }
            } catch (ConnectionException $e) {
                $this->connected = false;
                $nextReconnect = time() + 10;
            } catch (Exception $e) {
                if (!str_contains($e->getMessage(), 'no data received')) {
                    echo date('c') . " 读取异常: {$e->getMessage()}\n";
                }
            }

            if ($this->connected && time() - $this->lastHeartbeat >= 30) {
                $this->send(['code' => 10010, 'trace' => $this->trace()]);
                $this->lastHeartbeat = time();
            }
            usleep(100000);
        }
    }

    private function trace(): string
    {
        return bin2hex(random_bytes(16));
    }
}

$apiKey = getenv('INFOWAY_API_KEY') ?: 'yourApikey';
(new CryptoWebsocketClient($apiKey))->start();
```

{% endtab %}

{% tab title="JavaScript" %}

```javascript
class CryptoWebsocketClient {
  constructor(apiKey, business = "crypto") {
    this.wsUrl = `wss://data.infoway.io/ws?business=${encodeURIComponent(business)}&apikey=${encodeURIComponent(apiKey)}`;
    this.ws = null;
    this.heartbeatTimer = null;
    this.reconnectTimer = null;
    this.authFailed = false;
  }

  start() {
    this.connect();
  }

  connect() {
    if (this.authFailed) return;
    this.cleanupSocket();
    this.ws = new WebSocket(this.wsUrl);
    this.ws.onopen = () => {
      this.subscribe();
      this.startHeartbeat();
    };
    this.ws.onmessage = (event) => this.handle(String(event.data));
    this.ws.onerror = () => {};
    this.ws.onclose = (event) => {
      this.stopHeartbeat();
      if (event.code === 4001 || /401|507/.test(String(event.reason))) {
        this.authFailed = true;
        console.error("鉴权失败，停止重连");
        return;
      }
      this.scheduleReconnect();
    };
  }

  subscribe() {
    this.send({ code: 10000, trace: crypto.randomUUID(), data: { codes: "BTCUSDT" } });
    this.send({ code: 10003, trace: crypto.randomUUID(), data: { codes: "BTCUSDT" } });
    this.send({
      code: 10006,
      trace: crypto.randomUUID(),
      data: { arr: [{ type: 1, codes: "BTCUSDT" }] },
    });
  }

  send(msg) {
    if (!this.ws || this.ws.readyState !== WebSocket.OPEN) return;
    this.ws.send(JSON.stringify(msg));
  }

  handle(raw) {
    let msg;
    try {
      msg = JSON.parse(raw);
    } catch {
      console.log("忽略非 JSON 帧:", raw);
      return;
    }
    const { code, data, trace, msg: text } = msg;
    switch (code) {
      case 200:
        console.log("连接成功:", text);
        break;
      case 10002:
        console.log("成交推送:", data);
        break;
      case 10005:
        console.log("盘口推送:", data);
        break;
      case 10008:
        console.log("K线推送:", data);
        break;
      case 10001:
      case 10004:
      case 10007:
      case 11010:
        console.log("订阅确认", code, trace, text);
        break;
      case 10011:
        break;
      default:
        if (code >= 500 && code < 10000) {
          console.error("业务错误", code, text);
          if (code === 507) this.authFailed = true;
        } else {
          console.warn("未处理协议号", code, raw);
        }
    }
  }

  startHeartbeat() {
    this.stopHeartbeat();
    this.heartbeatTimer = setInterval(() => {
      this.send({ code: 10010, trace: crypto.randomUUID() });
    }, 30000);
  }

  stopHeartbeat() {
    if (this.heartbeatTimer) clearInterval(this.heartbeatTimer);
    this.heartbeatTimer = null;
  }

  scheduleReconnect() {
    if (this.authFailed || this.reconnectTimer) return;
    this.reconnectTimer = setTimeout(() => {
      this.reconnectTimer = null;
      this.connect();
    }, 10000);
  }

  cleanupSocket() {
    this.stopHeartbeat();
    if (this.ws) {
      this.ws.onclose = null;
      try { this.ws.close(); } catch {}
      this.ws = null;
    }
  }

  stop() {
    this.authFailed = true;
    if (this.reconnectTimer) clearTimeout(this.reconnectTimer);
    this.cleanupSocket();
  }
}

const apiKey = (typeof process !== "undefined" && process.env && process.env.INFOWAY_API_KEY) || "yourApikey";
const client = new CryptoWebsocketClient(apiKey, "crypto");
client.start();
```

{% endtab %}
{% endtabs %}
