You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 13:07:35