如何在Django中连接经纪商WebSocket并实时消费数据
在Django Channels中连接经纪商WebSocket并实时消费数据
现有代码的核心问题
你当前的实现里,self.external_socket.run_forever()是阻塞式调用,会占用Channels分配给Consumer的线程,导致Consumer无法响应前端的其他消息,甚至无法正常完成初始化流程。另外缺少经纪商WebSocket的认证步骤(Shoonya的WSTP连接后必须发送登录凭证才能订阅数据),也没有实现消息处理的具体逻辑。
修正后的实现方案
下面是调整后的完整代码,解决阻塞问题、添加认证流程,并实现数据转发:
1. 更新consumers.py
import json import threading import websocket from channels.generic.websocket import WebsocketConsumer class BrokerDataConsumer(WebsocketConsumer): def connect(self): # 先接受前端的WebSocket连接 self.accept() # 初始化外部经纪商WebSocket连接 self.external_socket = websocket.WebSocketApp( "wss://api.shoonya.com/NorenWSTP/", on_message=self.on_broker_message, on_error=self.on_broker_error, on_close=self.on_broker_close, on_open=self.on_broker_open # 添加连接成功后的回调(用于发送认证) ) # 启动独立线程运行外部WebSocket,避免阻塞Channels线程 self.external_thread = threading.Thread(target=self.external_socket.run_forever) self.external_thread.daemon = True # 设置为守护线程,随主线程退出 self.external_thread.start() self.send(text_data=json.dumps({ 'type': 'connection_established', 'message': '已连接到平台,正在对接经纪商...' })) def on_broker_open(self, ws): # 经纪商WebSocket连接成功后,发送认证信息(替换成你的实际凭证) auth_payload = json.dumps({ "t": "l", "uid": "你的经纪商用户ID", "pwd": "你的密码哈希", "factor2": "你的2FA验证码", "apkversion": "1.0.0", "source": "API" }) ws.send(auth_payload) self.send(text_data=json.dumps({ 'type': 'broker_auth_sent', 'message': '已向经纪商发送认证请求' })) def on_broker_message(self, ws, message): # 收到经纪商的数据后,转发给前端 self.send(text_data=json.dumps({ 'type': 'broker_data', 'data': json.loads(message) })) def on_broker_error(self, ws, error): # 处理经纪商WebSocket错误,通知前端 self.send(text_data=json.dumps({ 'type': 'broker_error', 'error': str(error) })) def on_broker_close(self, ws, close_status_code, close_msg): # 经纪商连接关闭时通知前端 self.send(text_data=json.dumps({ 'type': 'broker_disconnected', 'message': f'经纪商连接关闭: {close_msg}' })) def disconnect(self, close_code): # 当前端断开时,关闭经纪商WebSocket连接 if hasattr(self, 'external_socket'): self.external_socket.close() self.send(text_data=json.dumps({ 'type': 'disconnected', 'message': '已断开与平台的连接' }))
2. 更新routing.py(调整Consumer类名)
from django.urls import re_path from . import consumers websocket_urlpatterns = [ re_path(r'ws/broker-data/', consumers.BrokerDataConsumer.as_asgi()) ]
关键说明
- 线程隔离:把经纪商WebSocket的
run_forever()放到独立守护线程,避免阻塞Channels的事件循环,保证Consumer能正常处理前端的连接、断开等事件。 - 经纪商认证:Shoonya的WSTP协议要求连接成功后立即发送登录请求,
on_broker_open回调就是处理这个流程,需要替换成你自己的真实认证参数(注意密码需要是哈希值,具体参考Shoonya的API文档)。 - 数据转发:
on_broker_message把经纪商返回的实时行情/交易数据直接转发给前端,你可以在这里添加数据处理逻辑(比如过滤、转换格式)。 - 资源清理:在Consumer的
disconnect方法中关闭经纪商WebSocket,避免无效连接占用资源。
后续扩展建议
- 如果需要订阅特定品种的行情,可以在前端发送订阅指令后,通过Consumer调用
self.external_socket.send()向经纪商发送订阅请求。 - 可以添加日志记录,方便排查连接或数据问题。
- 若需要更高的性能,建议改用AsyncWebsocketConsumer,配合异步的WebSocket客户端(比如
websockets库),避免线程开销。
内容的提问来源于stack exchange,提问作者Shubham Agrawal
相关产品推荐
相关产品推荐

