Django Channels中channel_layer.group_send多次点击失效问题求助
Django Channels计时器功能异常问题排查与修复
问题描述
学习Django Channels时,实现基础计时器功能:点击按钮后,每秒通过WebSocket向React前端发送自点击时刻起的流逝秒数。但首次点击功能正常,后续点击时API请求成功响应,channel_layer.group_send却不再生效,last_hit时间未更新,计时器仍沿用初始点击的时间继续累加。
相关代码
consumer.py
import json import asyncio from channels.generic.websocket import AsyncWebsocketConsumer from django.utils import timezone class HelloConsumer(AsyncWebsocketConsumer): async def connect(self): # Join the hello_group when a WebSocket connection is established await self.channel_layer.group_add('hello_group', self.channel_name) await self.accept() async def disconnect(self, close_code): # Leave the hello_group when a WebSocket connection is closed await self.channel_layer.group_discard('hello_group', self.channel_name) async def receive(self, text_data): # Handle the received message if text_data == 'hello': # Update the last_hit timestamp self.scope['last_hit'] = timezone.now() else: await self.send(text_data=json.dumps({'error': 'Invalid message'})) async def send_hello_message(self, event): # Send hello message to the client while True: print('hehe2') last_hit = event.get('last_hit', None) if last_hit: seconds_ago = (timezone.now() - last_hit).seconds await self.send(text_data=f'Hello World hit {seconds_ago} seconds ago') else: await self.send(text_data='Hello World endpoint never hit lol') await asyncio.sleep(1) async def hello_event(self, event): print("reached here") # Handle the hello event from DRF view if event.get('hit', False): print('hehe') # await self.send_hello_message(event) print("last hit:", event.get('last_hit', None)) await asyncio.ensure_future(self.send_hello_message(event))
DRF视图
class HelloView(APIView): def get(self, request): # Trigger the WebSocket event when the DRF endpoint is hit channel_layer = get_channel_layer() async_to_sync(channel_layer.group_send)( 'hello_group', {'type': 'hello.event', 'hit': True, 'last_hit': timezone.now()} ) print("inside view") return Response({'message': 'Hello World!'})
routing.py
from channels.routing import URLRouter # from channels.auth import AuthMiddlewareStack from django.urls import path from .consumer import HelloConsumer url_router = URLRouter( [ path('ws/hello/', HelloConsumer.as_asgi()), ] )
asgi.py
""" ASGI config for channels_test project. It exposes the ASGI callable as a module-level variable named ``application``. For more information on this file, see https://docs.djangoproject.com/en/4.2/howto/deployment/asgi/ """ import os from django.core.asgi import get_asgi_application from channels.routing import ProtocolTypeRouter from channels.auth import AuthMiddlewareStack from channel_app.routing import url_router os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'channels_test.settings') # application = get_asgi_application() application = ProtocolTypeRouter({ 'http': get_asgi_application(), 'websocket': AuthMiddlewareStack( url_router ), })
React前端
import logo from './logo.svg'; import './App.css'; import React, { useState, useEffect } from 'react'; function App() { const [message, setMessage] = useState('Hello World endpoint never hit'); const hitHelloWorldEndpoint = async () => { try { // Trigger a request to the Hello World DRF endpoint const response = await fetch('http://localhost:8000/api/hello/'); // const data = await response.json(); const data = await response.json(); console.log('inside button', data.message) // setMessage(data.message); } catch (error) { console.error('Error hitting the Hello World endpoint:', error); } }; useEffect(() => { // Establish a WebSocket connection const socket = new WebSocket('ws://localhost:8000/ws/hello/'); // Handle WebSocket events socket.onopen = () => { console.log('WebSocket connection opened'); socket.send('hello'); }; socket.onmessage = (event) => { console.log('onmessage: ', event.data); try { const data = JSON.parse(event.data); setMessage(data); } catch (error) { console.error('Error parsing JSON:', error); if (typeof(event.data) === 'string'){ setMessage(event.data) } } }; socket.onclose = () => { console.log('WebSocket connection closed'); }; // Cleanup WebSocket connection on component unmount return () => { socket.close(); }; }, []); // Empty dependency array ensures useEffect runs only once on component mount return ( <div className="App"> <header className="App-header"> <img src={logo} className="App-logo" alt="logo" /> <p> {message} </p> <button onClick={hitHelloWorldEndpoint}>Hit Hello World Endpoint</button> </header> </div> ); } export default App;
运行日志
WebSocket CONNECT /ws/hello/ [127.0.0.1:52606] # this is the websocket connection successful message reached hereinside view #this is when the view is called and reached here means that it is inside the consumer class hehe last hit: 2023-12-14 19:04:24.685034+00:00 # this is called when the last_hit is updated hehe2 HTTP GET /api/hello/ 200 [0.00, 127.0.0.1:52607] #called when the api is hit hehe2 hehe2 hehe2 hehe2 hehe2 inside view HTTP GET /api/hello/ 200 [0.00, 127.0.0.1:52607] hehe2 hehe2 hehe2 hehe2 hehe2 inside view HTTP GET /api/hello/ 200 [0.00, 127.0.0.1:52607] hehe2 hehe2 hehe2 hehe2
问题原因分析
- Consumer被无限循环阻塞:
send_hello_message中的while True是无限循环,首次点击通过asyncio.ensure_future启动该任务后,会持续占用Consumer的事件循环,导致后续的hello_event事件无法被处理,所以第二次点击时,即使group_send发送了事件,Consumer也无法响应。 - 多任务冲突+状态未共享:每次点击会创建新的异步任务(如果没被阻塞的话),但旧任务仍在运行,且每个任务使用自己事件中的
last_hit值,无法共享更新后的时间;同时当前代码未维护Consumer级别的共享状态,导致新时间无法同步到运行中的计时器。
修复方案
修改后的consumer.py
import json import asyncio from channels.generic.websocket import AsyncWebsocketConsumer from django.utils import timezone class HelloConsumer(AsyncWebsocketConsumer): def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.last_hit = None self.timer_task = None # 跟踪计时器任务,避免多任务冲突 async def connect(self): await self.channel_layer.group_add('hello_group', self.channel_name) await self.accept() async def disconnect(self, close_code): await self.channel_layer.group_discard('hello_group', self.channel_name) # 断开连接时取消计时器任务,防止资源泄漏 if self.timer_task: self.timer_task.cancel() try: await self.timer_task except asyncio.CancelledError: pass async def receive(self, text_data): if text_data == 'hello': self.last_hit = timezone.now() else: await self.send(text_data=json.dumps({'error': 'Invalid message'})) async def send_hello_message(self): # 使用共享的self.last_hit,更新后立即生效 while True: print('hehe2') if self.last_hit: seconds_ago = (timezone.now() - self.last_hit).seconds await self.send(text_data=f'Hello World hit {seconds_ago} seconds ago') else: await self.send(text_data='Hello World endpoint never hit lol') await asyncio.sleep(1) async def hello_event(self, event): print("reached here") if event.get('hit', False): print('hehe') # 更新共享状态last_hit self.last_hit = event.get('last_hit', None) print("last hit:", self.last_hit) # 仅当计时器未运行时启动任务,避免重复创建 if not self.timer_task or self.timer_task.done(): self.timer_task = asyncio.create_task(self.send_hello_message())
修复关键点
- 共享状态维护:用
self.last_hit替代事件传递的last_hit,所有计时器逻辑共享同一状态,更新后下一次循环立即使用新时间。 - 任务跟踪与控制:通过
self.timer_task跟踪计时器任务,确保同一时间仅一个任务运行,避免多任务冲突;断开连接时取消任务,防止资源泄漏。 - 移除阻塞风险:虽然仍使用无限循环,但通过
asyncio.create_task管理任务,不会阻塞Consumer处理新的hello_event事件,后续点击能正常更新时间。
前端优化(可选)
后端发送的是纯字符串,无需JSON解析,可简化前端消息处理逻辑:
socket.onmessage = (event) => { console.log('onmessage: ', event.data); setMessage(event.data); };
内容的提问来源于stack exchange,提问作者lulnibba chad
相关产品推荐
相关产品推荐

