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
相关产品推荐
相关产品推荐

