连接方式
WebSocket 端点
wss://realtime-dex.chainstream.io/connection/websocket
连接认证
建立连接时需要在 URL 中提供 Access Token:wss://realtime-dex.chainstream.io/connection/websocket?token=YOUR_ACCESS_TOKEN
- 使用 SDK(推荐)
- 使用原生 WebSocket
- 命令行测试
SDK 已内置连接和认证处理,直接调用订阅方法即可:
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);
}
});
SDK 会自动检测连接状态,未连接时自动建立连接,无需手动调用
connect()。使用原生 WebSocket 时,连接后需要发送
connect 消息完成认证: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);
}
};
使用 连接后发送 connect 消息:
wscat 进行测试:wscat -c "wss://realtime-dex.chainstream.io/connection/websocket?token=YOUR_ACCESS_TOKEN"
{"connect":{"token":"YOUR_ACCESS_TOKEN","name":"test"},"id":1}
连接响应
认证成功后会收到如下响应:{
"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 线(支持 token / pool / pair 多种变体) |
| 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: | 单币 / 汇总钱包盈亏 |
| 排行榜 | 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: | 池子流动性快照 |
完整的订阅类型、参数说明和响应格式请参考 WebSocket API 文档。SDK(
client.stream.subscribeTokenCandles、subscribeTokenStats、subscribeTokenTrade 等)已经封装了这些频道字符串,无需手动拼接。订阅格式示例
// 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
}));
取消订阅
ws.send(JSON.stringify({
unsubscribe: {
channel: 'dex-candle:sol_6p6xgHyF7AeE6TZkSmFsko444wqoP15icUSqi2jfGiPN_1m'
},
id: 5
}));
消息格式
请求消息
Connect 消息(认证):{
"connect": {
"token": "YOUR_ACCESS_TOKEN",
"name": "client_name"
},
"id": 1
}
{
"subscribe": {
"channel": "dex-candle:sol_xxx_1m"
},
"id": 2
}
{
"unsubscribe": {
"channel": "dex-candle:sol_xxx_1m"
},
"id": 3
}
响应消息
订阅确认:{
"id": 2,
"subscribe": {}
}
{
"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
}
}
}
}
{
"id": 2,
"error": {
"code": 100,
"message": "invalid channel"
}
}
心跳保活
WebSocket 连接需要定期发送心跳消息以保持活跃。根据 connect 响应中的ping 字段(通常为 25 秒),在此间隔内发送心跳:
// 心跳消息
ws.send(JSON.stringify({}));
// 或者发送 ping
ws.send(JSON.stringify({ ping: {} }));
如果在指定时间内(通常为 ping 间隔的 3 倍)未发送任何消息,服务器将主动断开连接。
完整示例
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);
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())
最佳实践
性能优化
使用过滤条件
只订阅需要的数据,减少带宽消耗。使用 CEL 表达式过滤数据。
批量处理
对高频数据进行批量处理而非逐条处理,使用消息队列缓冲。
本地缓存
缓存 Token 信息等静态数据,减少重复处理。
连接复用
单个连接可订阅多个频道,避免创建多个连接。
错误处理
- 监听错误事件 — 及时处理连接错误和数据错误
- 实现重试机制 — 使用指数退避策略进行重连
- 日志记录 — 记录关键事件便于问题排查
- 优雅降级 — WebSocket 不可用时切换到轮询
资源管理
// ✅ 及时取消不需要的订阅
client.unsubscribe('dex-candle:sol_xxx_1m');
// ✅ 优雅关闭连接
function gracefulClose() {
// 1. 停止心跳
client.stopPing();
// 2. 关闭连接
client.close();
}
相关文档
WebSocket API 参考
完整的订阅类型和参数说明
价格预警机器人
实战:构建价格监控 Bot

