> ## Documentation Index
> Fetch the complete documentation index at: https://docs.chainstream.io/llms.txt
> Use this file to discover all available pages before exploring further.

# 實時資料流

> WebSocket 實時資料流接入指南

ChainStream 提供強大的實時資料流處理能力，讓開發者能夠即時接收鏈上事件、交易和狀態變化。本文件介紹 WebSocket 連線、訂閱機制和最佳實踐。

***

## 連線方式

### WebSocket 端點

```
wss://realtime-dex.chainstream.io/connection/websocket
```

### 連線認證

建立連線時需要在 URL 中提供 Access Token：

```
wss://realtime-dex.chainstream.io/connection/websocket?token=YOUR_ACCESS_TOKEN
```

<Tabs>
  <Tab title="使用 SDK（推薦）">
    SDK 已內建連線和認證處理，直接呼叫訂閱方法即可：

    ```javascript theme={null}
    import { ChainStreamClient } from '@chainstream-io/sdk';
    import { Resolution } from '@chainstream-io/sdk/openapi';

    const client = new ChainStreamClient(process.env.CHAINSTREAM_ACCESS_TOKEN);

    // 直接订阅，SDK 自动处理连接和认证
    client.stream.subscribeTokenCandles({
      chain: 'sol',
      tokenAddress: '6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN',
      resolution: Resolution._1m,
      callback: (data) => {
        console.log('收到数据:', data);
      }
    });
    ```

    <Note>
      SDK 會自動檢測連線狀態，未連線時自動建立連線，無需手動呼叫 `connect()`。
    </Note>
  </Tab>

  <Tab title="使用原生 WebSocket">
    使用原生 WebSocket 時，連線後需要傳送 `connect` 訊息完成認證：

    ```javascript theme={null}
    const token = process.env.CHAINSTREAM_ACCESS_TOKEN;
    const ws = new WebSocket(
      `wss://realtime-dex.chainstream.io/connection/websocket?token=${token}`
    );

    ws.onopen = () => {
      console.log('WebSocket 连接已建立');
      
      // 发送 connect 消息完成认证
      ws.send(JSON.stringify({
        connect: {
          token: token,
          name: 'js'
        },
        id: 1
      }));
    };

    ws.onmessage = (event) => {
      const data = JSON.parse(event.data);
      
      // 处理 connect 响应
      if (data.connect) {
        console.log('✅ 认证成功，客户端 ID:', data.connect.client);
        
        // 认证成功后开始订阅
        ws.send(JSON.stringify({
          subscribe: {
            channel: 'dex-candle:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN_1m'
          },
          id: 2
        }));
      }
      
      // 处理订阅数据
      if (data.push) {
        console.log('收到数据:', data.push.pub.data);
      }
    };
    ```
  </Tab>

  <Tab title="命令列測試">
    使用 `wscat` 進行測試：

    ```bash theme={null}
    wscat -c "wss://realtime-dex.chainstream.io/connection/websocket?token=YOUR_ACCESS_TOKEN"
    ```

    連線後傳送 connect 訊息：

    ```json theme={null}
    {"connect":{"token":"YOUR_ACCESS_TOKEN","name":"test"},"id":1}
    ```
  </Tab>
</Tabs>

### 連線響應

認證成功後會收到如下響應：

```json theme={null}
{
  "id": 1,
  "connect": {
    "client": "0f819f5f-7d8b-4949-9433-0e91bbfe1cdb",
    "version": "0.0.0 OSS",
    "expires": true,
    "ttl": 86002,
    "ping": 25,
    "pong": true
  }
}
```

| 欄位       | 說明              |
| :------- | :-------------- |
| `client` | 客戶端唯一標識         |
| `ttl`    | Token 剩餘有效時間（秒） |
| `ping`   | 心跳間隔（秒）         |
| `pong`   | 是否支援 pong 響應    |

***

## 訂閱型別

ChainStream WebSocket 支援多種資料訂閱型別：

| 類別         | 訂閱頻道字首                                                                                                                                                                  | 說明                     |
| :--------- | :---------------------------------------------------------------------------------------------------------------------------------------------------------------------- | :--------------------- |
| K 線（USD）   | `dex-candle:`                                                                                                                                                           | 代幣 K 線（代幣 / 池 / 交易對變體） |
| K 線（原生幣計價） | `dex-candle-in-native:`                                                                                                                                                 | 同上，以該鏈原生資產計價           |
| 池 K 線      | `dex-pool-candle:` / `dex-pair-candle:`                                                                                                                                 | 池級或交易對級 K 線            |
| 代幣統計       | `dex-token-stats:`                                                                                                                                                      | 多視窗成交統計（1m、5m…1W、1M）   |
| 持有者統計      | `dex-token-holding:`                                                                                                                                                    | 持幣分佈與餘額標籤              |
| 代幣供給       | `dex-token-supply:`                                                                                                                                                     | 供給與市值更新                |
| 流動性        | `dex-token-liquidity:` / `dex-token-total-liquidity:`                                                                                                                   | 最大池 / 總流動性             |
| 新代幣        | `dex-new-token:` / `dex-new-tokens:` / `dex-new-tokens-metadata:`                                                                                                       | 新上線（單條 / 批次 / 後設資料）    |
| 代幣交易       | `dex-trade:`                                                                                                                                                            | 按代幣地址過濾的成交             |
| 錢包餘額       | `dex-wallet-balance:`                                                                                                                                                   | 錢包代幣餘額變動               |
| 錢包交易       | `dex-wallet-trade:`                                                                                                                                                     | 按錢包地址過濾的成交             |
| 錢包 PnL     | `dex-wallet-token-pnl:` / `dex-wallet-pnl-list:`                                                                                                                        | 逐幣 / 聚合錢包 PnL          |
| 排名         | `dex-ranking-list:` / `dex-ranking-token-stats-list:` / `dex-ranking-token-holding-list:` / `dex-ranking-token-supply-list:` / `dex-ranking-token-bounding-curve-list:` | 榜單成員 + 逐幣統計            |
| DEX 池      | `dex-pool-balance:`                                                                                                                                                     | 池流動性快照                 |

<Tip>
  完整的訂閱型別、引數說明和響應格式請參考 [WebSocket API 文件](/zh-Hant/api-reference/endpoint/websocket/api)。SDK（`client.stream.subscribeTokenCandles`、`subscribeTokenStats`、`subscribeTokenTrade`……）已封裝好這些頻道字串，不必手工拼接。
</Tip>

### 訂閱格式示例

```javascript theme={null}
// K线数据订阅
ws.send(JSON.stringify({
  subscribe: {
    channel: 'dex-candle:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN_1m'
  },
  id: 2
}));

// 代币统计订阅
ws.send(JSON.stringify({
  subscribe: {
    channel: 'dex-token-stats:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN'
  },
  id: 3
}));

// 新代币订阅
ws.send(JSON.stringify({
  subscribe: {
    channel: 'dex-new-token:sol'
  },
  id: 4
}));
```

### 取消訂閱

```javascript theme={null}
ws.send(JSON.stringify({
  unsubscribe: {
    channel: 'dex-candle:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN_1m'
  },
  id: 5
}));
```

***

## 訊息格式

### 請求訊息

**Connect 訊息（認證）：**

```json theme={null}
{
  "connect": {
    "token": "YOUR_ACCESS_TOKEN",
    "name": "client_name"
  },
  "id": 1
}
```

**Subscribe 訊息（訂閱）：**

```json theme={null}
{
  "subscribe": {
    "channel": "dex-candle:sol_xxx_1m"
  },
  "id": 2
}
```

**Unsubscribe 訊息（取消訂閱）：**

```json theme={null}
{
  "unsubscribe": {
    "channel": "dex-candle:sol_xxx_1m"
  },
  "id": 3
}
```

### 響應訊息

**訂閱確認：**

```json theme={null}
{
  "id": 2,
  "subscribe": {}
}
```

**資料推送：**

```json theme={null}
{
  "push": {
    "channel": "dex-candle:sol_xxx_1m",
    "pub": {
      "data": {
        "o": 0.001234,
        "c": 0.001256,
        "h": 0.001280,
        "l": 0.001200,
        "v": 1234567,
        "t": 1706745600
      }
    }
  }
}
```

**錯誤訊息：**

```json theme={null}
{
  "id": 2,
  "error": {
    "code": 100,
    "message": "invalid channel"
  }
}
```

***

## 心跳保活

WebSocket 連線需要定期傳送心跳訊息以保持活躍。根據 connect 響應中的 `ping` 欄位（通常為 25 秒），在此間隔內傳送心跳：

```javascript theme={null}
// 心跳消息
ws.send(JSON.stringify({}));

// 或者发送 ping
ws.send(JSON.stringify({ ping: {} }));
```

<Warning>
  如果在指定時間內（通常為 ping 間隔的 3 倍）未傳送任何訊息，伺服器將主動斷開連線。
</Warning>

***

## 完整示例

<CodeGroup>
  ```javascript JavaScript theme={null}
  const WebSocket = require('ws');

  class ChainStreamWebSocket {
    constructor(accessToken) {
      this.accessToken = accessToken;
      this.ws = null;
      this.messageId = 0;
      this.reconnectAttempts = 0;
      this.maxReconnectAttempts = 10;
      this.subscriptions = new Set();
      this.pingInterval = null;
    }

    connect() {
      const url = `wss://realtime-dex.chainstream.io/connection/websocket?token=${this.accessToken}`;
      this.ws = new WebSocket(url);

      this.ws.onopen = () => {
        console.log('WebSocket 连接已建立');
        this.reconnectAttempts = 0;
        
        // 发送 connect 消息
        this.send({
          connect: {
            token: this.accessToken,
            name: 'nodejs'
          }
        });
      };

      this.ws.onmessage = (event) => {
        const data = JSON.parse(event.data);
        this.handleMessage(data);
      };

      this.ws.onclose = (event) => {
        console.log(`连接关闭: ${event.code}`);
        this.stopPing();
        
        if (event.code !== 1000) {
          this.scheduleReconnect();
        }
      };

      this.ws.onerror = (error) => {
        console.error('WebSocket 错误:', error.message);
      };
    }

    handleMessage(data) {
      // 处理 connect 响应
      if (data.connect) {
        console.log('✅ 认证成功');
        this.startPing(data.connect.ping || 25);
        this.resubscribe();
        return;
      }

      // 处理数据推送
      if (data.push) {
        console.log(`[${data.push.channel}]`, data.push.pub.data);
        return;
      }

      // 处理错误
      if (data.error) {
        console.error('错误:', data.error.message);
        return;
      }
    }

    send(message) {
      message.id = ++this.messageId;
      this.ws.send(JSON.stringify(message));
    }

    subscribe(channel) {
      this.subscriptions.add(channel);
      this.send({ subscribe: { channel } });
    }

    unsubscribe(channel) {
      this.subscriptions.delete(channel);
      this.send({ unsubscribe: { channel } });
    }

    resubscribe() {
      this.subscriptions.forEach(channel => {
        this.send({ subscribe: { channel } });
      });
    }

    startPing(interval) {
      this.pingInterval = setInterval(() => {
        if (this.ws.readyState === WebSocket.OPEN) {
          this.ws.send('{}');
        }
      }, interval * 1000);
    }

    stopPing() {
      if (this.pingInterval) {
        clearInterval(this.pingInterval);
        this.pingInterval = null;
      }
    }

    scheduleReconnect() {
      if (this.reconnectAttempts >= this.maxReconnectAttempts) {
        console.error('❌ 达到最大重连次数');
        return;
      }

      const delay = Math.min(1000 * Math.pow(2, this.reconnectAttempts), 30000);
      console.log(`⏳ ${delay}ms 后重连... (第 ${this.reconnectAttempts + 1} 次)`);
      this.reconnectAttempts++;

      setTimeout(() => this.connect(), delay);
    }

    close() {
      this.stopPing();
      if (this.ws) {
        this.ws.close(1000, 'Normal closure');
      }
    }
  }

  // 使用示例
  const client = new ChainStreamWebSocket(process.env.CHAINSTREAM_ACCESS_TOKEN);
  client.connect();

  // 等待连接成功后订阅
  setTimeout(() => {
    client.subscribe('dex-candle:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN_1m');
    client.subscribe('dex-new-token:sol');
  }, 1000);
  ```

  ```python Python theme={null}
  import asyncio
  import websockets
  import json
  import os

  class ChainStreamWebSocket:
      def __init__(self, access_token):
          self.access_token = access_token
          self.ws = None
          self.message_id = 0
          self.subscriptions = set()
          self.ping_interval = 25
      
      async def connect(self):
          url = f"wss://realtime-dex.chainstream.io/connection/websocket?token={self.access_token}"
          
          async with websockets.connect(url) as ws:
              self.ws = ws
              print("WebSocket 连接已建立")
              
              # 发送 connect 消息
              await self.send({
                  "connect": {
                      "token": self.access_token,
                      "name": "python"
                  }
              })
              
              # 启动心跳和消息接收
              await asyncio.gather(
                  self.heartbeat(),
                  self.receive()
              )
      
      async def send(self, message):
          self.message_id += 1
          message["id"] = self.message_id
          await self.ws.send(json.dumps(message))
      
      async def receive(self):
          async for message in self.ws:
              data = json.loads(message)
              await self.handle_message(data)
      
      async def handle_message(self, data):
          if "connect" in data:
              print("✅ 认证成功")
              self.ping_interval = data["connect"].get("ping", 25)
              await self.resubscribe()
          elif "push" in data:
              channel = data["push"]["channel"]
              pub_data = data["push"]["pub"]["data"]
              print(f"[{channel}]", pub_data)
          elif "error" in data:
              print(f"错误: {data['error']['message']}")
      
      async def subscribe(self, channel):
          self.subscriptions.add(channel)
          await self.send({"subscribe": {"channel": channel}})
      
      async def resubscribe(self):
          for channel in self.subscriptions:
              await self.send({"subscribe": {"channel": channel}})
      
      async def heartbeat(self):
          while True:
              await asyncio.sleep(self.ping_interval)
              if self.ws and self.ws.open:
                  await self.ws.send("{}")

  # 使用示例
  async def main():
      client = ChainStreamWebSocket(os.environ["CHAINSTREAM_ACCESS_TOKEN"])
      
      # 预先添加订阅
      client.subscriptions.add("dex-candle:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN_1m")
      client.subscriptions.add("dex-new-token:sol")
      
      await client.connect()

  asyncio.run(main())
  ```
</CodeGroup>

***

## 最佳實踐

### 效能最佳化

<CardGroup cols={2}>
  <Card title="使用過濾條件" icon="filter">
    只訂閱需要的資料，減少頻寬消耗。使用 CEL 表示式過濾資料。
  </Card>

  <Card title="批次處理" icon="layer-group">
    對高頻資料進行批次處理而非逐條處理，使用訊息佇列緩衝。
  </Card>

  <Card title="本地快取" icon="database">
    快取 Token 資訊等靜態資料，減少重複處理。
  </Card>

  <Card title="連線複用" icon="plug">
    單個連線可訂閱多個頻道，避免建立多個連線。
  </Card>
</CardGroup>

### 錯誤處理

1. **監聽錯誤事件** — 及時處理連線錯誤和資料錯誤
2. **實現重試機制** — 使用指數退避策略進行重連
3. **日誌記錄** — 記錄關鍵事件便於問題排查
4. **優雅降級** — WebSocket 不可用時切換到輪詢

### 資源管理

```javascript theme={null}
// ✅ 及时取消不需要的订阅
client.unsubscribe('dex-candle:sol_xxx_1m');

// ✅ 优雅关闭连接
function gracefulClose() {
  // 1. 停止心跳
  client.stopPing();
  
  // 2. 关闭连接
  client.close();
}
```

***

## 相關文件

<CardGroup cols={2}>
  <Card title="WebSocket API 參考" icon="plug" href="/zh-Hant/api-reference/endpoint/websocket/api">
    完整的訂閱型別和引數說明
  </Card>

  <Card title="價格預警機器人" icon="bell" href="/zh-Hant/docs/tutorials/build-price-alert-bot">
    實戰：構建價格監控 Bot
  </Card>
</CardGroup>
