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

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

问题原因分析

  1. Consumer被无限循环阻塞:send_hello_message中的while True是无限循环,首次点击通过asyncio.ensure_future启动该任务后,会持续占用Consumer的事件循环,导致后续的hello_event事件无法被处理,所以第二次点击时,即使group_send发送了事件,Consumer也无法响应。
  2. 多任务冲突+状态未共享:每次点击会创建新的异步任务(如果没被阻塞的话),但旧任务仍在运行,且每个任务使用自己事件中的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 00:57:33