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

Django Channels Websocket连接异常:刷新后无法重连,进程被终止

环境信息

  • Ubuntu 16.04.6
  • conda 4.12.0
  • Apache/2.4.18 (Ubuntu)
  • python==3.8.1
  • Django==4.0.3
  • channels==3.0.5
  • asgi-redis==1.4.3
  • asgiref==3.4.1
  • daphne==3.0.2

我正在搭建一个仅向已认证用户转发Redis消息的WebSocket服务,用户之间无需通信,因此不需要Channel Layers(据我了解这是Channels的可选组件)。

我通过自定义中间件完成用户认证,从会话ID中获取已认证用户信息及user_id并附加到scope中,仅向特定用户广播消息。我选择通过Apache2的ProxyPass将所有流量路由到8033端口的Daphne,实现单域名部署。

但现在难以维持WebSocket连接,尤其是刷新浏览器时:首次请求可正常工作,后续请求失败,journalctl中出现如下报错:

Application instance <Task pending name='Task-22' coro=<ProtocolTypeRouter.__call__() running at /root/miniconda2/lib/python3.8/site-packages/channels/routing.py:71> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7f658c03f670>()]>> for connection <WebSocketProtocol client=['127.0.0.1', 46010] path=b'/ws/userupdates/'> took too long to shut down and was killed.

我在Channels的GitHub仓库里找了大量解决方案,但仍未解决问题。当前代码能建立初始连接、返回连接响应{"success": true, "user_id": XXXXXX, "message": "Connected"}并正常转发Redis消息,但刷新浏览器或关闭重开后无法建立连接,必须重启Apache和Daphne才可恢复。

我猜测问题出在未正确断开消费者或async await使用不当,请问有什么解决思路?


相关Apache配置

RewriteEngine On
RewriteCond %{HTTP:Connection} Upgrade [NC]
RewriteCond %{HTTP:Upgrade} websocket [NC]
RewriteRule /(.*) ws://127.0.0.1:8033/$1 [P,L]

<Location />
    ProxyPass http://127.0.0.1:8033/
    ProxyPassReverse /
</Location>

app/settings.py

[...]
ASGI_APPLICATION = 'app.asgi.application'
ASGI_THREADS = 1000
CHANNEL_LAYERS = {}
[...]

app/asgi.py

import os

from django.core.asgi import get_asgi_application
from django.urls import path
from channels.routing import ProtocolTypeRouter, URLRouter
from channels.security.websocket import AllowedHostsOriginValidator

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'app.settings')

import websockets.routing
from user.models import UserAuthMiddleware

application = ProtocolTypeRouter({
    'http': get_asgi_application(),
    "websocket": AllowedHostsOriginValidator(
        UserAuthMiddleware(
            URLRouter(websockets.routing.websocket_urlpatterns)
        )
    ),
})

user/models.py::UserAuthMiddleware

说明:从自定义认证层获取user_id并附加到scope中

class UserAuthMiddleware(CookieMiddleware):
    def __init__(self, app):
        self.app = app
    async def __call__(self, scope, receive, send):
        # Check this actually has headers. They're a required scope key for HTTP and WS.
        if "headers" not in scope:
            raise UserSessionError(
                "UserAuthMiddleware was passed a scope that did not have a headers key "
                + "(make sure it is only passed HTTP or WebSocket connections)"
            )
        # Go through headers to find the cookie one
        for name, value in scope.get("headers", []):
            if name == b"cookie":
                cookies = parse_cookie(value.decode("latin1"))
                break
        else:
            # No cookie header found - add an empty default.
            cookies = {}

        # now gather user data from session
        try:
            req = HttpRequest()
            req.GET = QueryDict(query_string=scope.get("query_string"))
            setattr(req, 'COOKIES', cookies)
            setattr(req, 'headers', scope.get("headers")),

            session = UserSession(req)
            scope['user_id'] = session.get_user_id()
        except UserSessionError as e:
            raise e

        return await self.app(scope, receive, send)

websockets/routing.py

from django.urls import re_path
from . import consumers

websocket_urlpatterns = [
    re_path(r'ws/userupdates/', consumers.UserUpdatesConsumer.as_asgi())
]

websockets/consumers.py::UserUpdatesConsumer

from channels.generic.websocket import JsonWebsocketConsumer
import json, redis

class UserUpdatesConsumer(JsonWebsocketConsumer):
    def connect(self):
        self.accept()

        self.redis = redis.Redis(host='127.0.0.1', port=6379, db=0, decode_responses=True)
        self.p = self.redis.pubsub()

        if 'user_id' not in self.scope:
            self.send_json({
                'success': False,
                'message': 'No user_id present'
            })
            self.close()
        else:
            self.send_json({
                'success': True,
                'user_id': self.scope['user_id'],
                'message': 'Connected'
            })

            self.p.psubscribe(f"dip_alerts")
            self.p.psubscribe(f"userupdates_{self.scope['user_id']}*")
            for message in self.p.listen():

                if message.get('type') == 'psubscribe' and message.get('data') in [1,2]:
                    continue

                if message.get('channel') == "dip_alerts":
                    self.send_json({
                        "key": "dip_alerts",
                        "event": "dip_alert",
                        "data": json.loads(message.get('data'))
                    })
                else:
                    self.send(message.get('data'))

内容的提问来源于stack exchange,提问作者Mike Christensen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 20:01:07