如何为Python Websocket连接指定多源IP以降低KuCoin行情延迟?
问题
我用Python的websocket库开发了KuCoin全市场Level2行情抓取程序,通过多线程订阅约1200个实时频道,每个线程订阅一组交易对(比如Thread1订阅BTC-USDT、ETC-USDT等,Thread2订阅ETH-USDT、ETH-BTC等)。
服务器配置为4核Xeon、4GB内存、10Gbps端口、8个可用IP。目前程序默认用系统单IP连接KuCoin的Websocket服务器,因API限制导致大量订阅出现数据延迟,希望给每个Websocket连接分配独立的源IP,实现多源IP分散请求。
现有核心代码如下:
import websocket, kucoin import json from threading import Thread def func_do_some_works_on_data(message): # 处理行情数据的逻辑 pass def proc_websocket(symbols): def on_ping(wsapp, message):pass def on_pong(wsapp, message):pass def on_error(wsapp, message):print(message) def on_open(wsapp): wsapp.send(json.dumps({'type':'subscribe', 'topic':'/market/level2:' + ','.join(symbols), 'response':True})) def on_message(wsapp, message): func_do_some_works_on_data(message) wss_endpoint = kucoin.get_ws_endpoint() ws_endpoint_url = f"{wss_endpoint['instanceServers'][0]['endpoint']}?token={wss_endpoint['token']}" while True: wsapp = websocket.WebSocketApp( ws_endpoint_url, on_message=on_message, on_open=on_open, on_error=on_error, on_ping=on_ping, on_pong=on_pong ) wsapp.run_forever(ping_interval=15, ping_timeout=10) symbols = kucoin.get_symbols_all() Thread(target=proc_websocket, args=(symbols[:100],)).start() Thread(target=proc_websocket, args=(symbols[100:200],)).start() Thread(target=proc_websocket, args=(symbols[200:300],)).start()
解决方案
方法1:基于现有websocket库绑定源IP
websocket.WebSocketApp的run_forever方法支持传入自定义socket,你可以通过绑定指定源IP的socket,实现每个线程用独立IP连接。
修改后代码示例:
import websocket, kucoin import json import socket from threading import Thread # 服务器上的8个可用源IP列表,替换为你的实际IP SOURCE_IPS = ["192.168.1.101", "192.168.1.102", "192.168.1.103", "192.168.1.104", "192.168.1.105", "192.168.1.106", "192.168.1.107", "192.168.1.108"] def func_do_some_works_on_data(message): # 处理行情数据的逻辑 pass def proc_websocket(symbols, source_ip): def on_ping(wsapp, message):pass def on_pong(wsapp, message):pass def on_error(wsapp, message):print(message) def on_open(wsapp): wsapp.send(json.dumps({'type':'subscribe', 'topic':'/market/level2:' + ','.join(symbols), 'response':True})) def on_message(wsapp, message): func_do_some_works_on_data(message) wss_endpoint = kucoin.get_ws_endpoint() ws_endpoint_url = f"{wss_endpoint['instanceServers'][0]['endpoint']}?token={wss_endpoint['token']}" # 创建绑定指定源IP的socket def create_bound_socket(): sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.bind((source_ip, 0)) # 0表示自动选择可用端口 return sock while True: wsapp = websocket.WebSocketApp( ws_endpoint_url, on_message=on_message, on_open=on_open, on_error=on_error, on_ping=on_ping, on_pong=on_pong ) # 传入绑定好源IP的socket wsapp.run_forever( ping_interval=15, ping_timeout=10, socket=create_bound_socket() ) symbols = kucoin.get_symbols_all() # 按8个IP平均拆分交易对分组 symbol_groups = [symbols[i::8] for i in range(8)] # 启动线程,每个线程对应一个源IP for idx, group in enumerate(symbol_groups): if group: # 跳过空分组 Thread(target=proc_websocket, args=(group, SOURCE_IPS[idx])).start()
方法2:改用异步websockets库(性能更优)
如果追求更高的资源利用率,推荐使用异步的websockets库(注意与你当前使用的websocket库不是同一个),同样支持绑定源IP:
import asyncio import websockets import socket import json from kucoin import get_ws_endpoint, get_symbols_all # 替换为你的实际源IP列表 SOURCE_IPS = ["192.168.1.101", "192.168.1.102", ..., "192.168.1.108"] def func_do_some_works_on_data(message): # 处理行情数据的逻辑 pass async def connect_with_source_ip(symbols, source_ip): wss_endpoint = get_ws_endpoint() ws_url = f"{wss_endpoint['instanceServers'][0]['endpoint']}?token={wss_endpoint['token']}" # 创建绑定源IP的socket sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.bind((source_ip, 0)) # 解析KuCoin Websocket目标地址 target_host = ws_url.replace("wss://", "").split("?")[0].split(":")[0] target_port = 443 # 建立TCP连接 await asyncio.get_event_loop().sock_connect(sock, (target_host, target_port)) # 升级为Websocket连接并订阅行情 async with websockets.connect( ws_url, sock=sock, ping_interval=15, ping_timeout=10 ) as ws: await ws.send(json.dumps({ 'type':'subscribe', 'topic':'/market/level2:' + ','.join(symbols), 'response':True })) # 持续接收数据 async for message in ws: func_do_some_works_on_data(message) async def main(): symbols = get_symbols_all() symbol_groups = [symbols[i::8] for i in range(8)] tasks = [] for idx, group in enumerate(symbol_groups): if group: tasks.append(connect_with_source_ip(group, SOURCE_IPS[idx])) await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
注意事项
- 确保所有源IP都能正常访问KuCoin的Websocket服务器,无防火墙或路由限制。
- 拆分交易对时尽量保证每个分组数量均衡,避免单个IP负载过高。
- 可添加重连容错逻辑,避免单个连接断开影响整个分组的订阅。
内容的提问来源于stack exchange,提问作者MPoorzahed
相关产品推荐
相关产品推荐

