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

Django Channels Consumer无法接收模型保存时发送的新通知

问题描述

我编写了一个WebSocket Consumer,用于获取所有未读通知,以及客户端保持WebSocket连接期间新增的通知。目前可以成功连接Consumer并查看已有通知,但通过模型save方法创建并发送新通知时,客户端无法接收到这些新通知。

consumers.py

class NotificationConsumer(WebsocketConsumer):
    def get_notifications(self, user):
        new_followers = NotificationSerializer(Notification.objects.filter(content=1), many=True)

        notifs = {
            "new_followers": new_followers.data
        }

        return {
            "count": sum(len(notif) for notif in notifs.values()),
            "notifs": notifs
        }

    def connect(self):
        user = self.scope["user"]

        if user:
            self.channel_name = user.username
            self.room_group_name = 'notification_%s' % self.channel_name

            notification = self.get_notifications(self.scope["user"])
            print("Notification", notification)

            self.accept()
            return self.send(text_data=json.dumps(notification))
        
        self.disconnect(401)
    
    async def disconnect(self, close_code):
        await self.channel_layer.group_discard(
            self.room_group_name,
            self.channel_name
        )
    
    async def receive(self, text_data):
        text_data_json = json.loads(text_data)
        count = text_data_json['count']
        notification_type = text_data_json['notification_type']
        notification = text_data_json['notification']
        print(text_data_json)

        await self.channel_layer.group_send(
            self.room_group_name,
            {
                "type": "receive",
                "count": count,
                "notification_type": notification_type,
                "notification": notification
            }
        )

models.py

class Notification(models.Model):
    content = models.CharField(max_length=200, choices=NOTIFICATION_CHOICES, default=1)
    additional_info = models.CharField(max_length=100, blank=True)
    seen = models.BooleanField(default=False)
    users_notified = models.ManyToManyField("core.User", blank=True)
    user_signaling = models.ForeignKey("core.User", on_delete=models.CASCADE, related_name="user_signaling")

    def save(self, *args, **kwargs):
        if self._state.adding:
            super(Notification, self).save(*args, **kwargs)

        else:
            notif = Notification.objects.filter(seen=False)
            data = {
                "type": "receive",
                "count": len(notif),
                "notification_type": self.content,
                "notification": self.additional_info
            }

            channel_layer = get_channel_layer()
            for user in self.users_notified.all():
                async_to_sync(channel_layer.group_send(
                        f"notification_{user.username}",
                        data
                    )
                )
                print(channel_layer)

示例用例

notif = Notification.objects.create(
      content=1, 
      user_signaling=user,
      additional_info=f"{user.username} just followed you!"
)
notif.users_notified.add(other_user)
notif.save()

我已通过检查self._state.adding,确保仅在模型创建后添加users_notified时才发送WebSocket消息。


问题分析与修复

核心问题

  1. Consumer未加入分组:connect方法仅定义了分组名称,但未调用group_add将当前连接加入对应通知分组,导致模型发送的分组消息无法被Consumer接收。
  2. 异步方法不匹配:connect是同步方法,而disconnect、receive是异步方法,Django Channels要求Consumer方法保持统一的同步/异步模式,否则会出现上下文异常。
  3. 分组消息处理缺失:发送的消息type指定为receive,但该方法是处理客户端发来的消息,并非处理分组广播的消息,缺少专门的分组消息处理方法。

修复步骤

1. 修正consumers.py,统一异步模式并加入分组

将connect改为异步方法,添加分组加入逻辑,新增分组消息处理方法:

import json
from channels.generic.websocket import AsyncWebsocketConsumer
from .serializers import NotificationSerializer
from .models import Notification
from asgiref.sync import sync_to_async

class NotificationConsumer(AsyncWebsocketConsumer):
    async def get_notifications(self, user):
        def get_unread_notifs():
            new_followers = NotificationSerializer(Notification.objects.filter(content=1), many=True)
            notifs = {
                "new_followers": new_followers.data
            }
            return {
                "count": sum(len(notif) for notif in notifs.values()),
                "notifs": notifs
            }
        
        return await sync_to_async(get_unread_notifs)()

    async def connect(self):
        user = self.scope["user"]

        if user.is_authenticated:
            self.room_group_name = f'notification_{user.username}'

            # 将当前连接加入分组
            await self.channel_layer.group_add(
                self.room_group_name,
                self.channel_name
            )

            notification = await self.get_notifications(user)
            print("Notification", notification)

            await self.accept()
            await self.send(text_data=json.dumps(notification))
            return
        
        await self.disconnect(401)
    
    async def disconnect(self, close_code):
        # 从分组移除当前连接
        await self.channel_layer.group_discard(
            self.room_group_name,
            self.channel_name
        )
    
    # 处理客户端发送的消息(无需可删除)
    async def receive(self, text_data):
        text_data_json = json.loads(text_data)
        print(text_data_json)
    
    # 新增分组消息处理方法,对应模型发送的type字段
    async def notification_message(self, event):
        data = {
            "count": event["count"],
            "notification_type": event["notification_type"],
            "notification": event["notification"]
        }
        await self.send(text_data=json.dumps(data))

2. 修正models.py中的消息type字段

将发送的消息type改为notification_message,确保能被Consumer正确处理:

from django.db import models
from channels.layers import get_channel_layer
from asgiref.sync import async_to_sync

class Notification(models.Model):
    content = models.CharField(max_length=200, choices=NOTIFICATION_CHOICES, default=1)
    additional_info = models.CharField(max_length=100, blank=True)
    seen = models.BooleanField(default=False)
    users_notified = models.ManyToManyField("core.User", blank=True)
    user_signaling = models.ForeignKey("core.User", on_delete=models.CASCADE, related_name="user_signaling")

    def save(self, *args, **kwargs):
        is_new = self._state.adding
        super().save(*args, **kwargs)

        if not is_new:
            # 计算当前用户的未读通知数
            def get_unread_count(user):
                return Notification.objects.filter(users_notified=user, seen=False).count()
            
            channel_layer = get_channel_layer()
            for user in self.users_notified.all():
                unread_count = get_unread_count(user)
                data = {
                    "type": "notification_message",
                    "count": unread_count,
                    "notification_type": self.content,
                    "notification": self.additional_info
                }
                async_to_sync(channel_layer.group_send)(
                    f"notification_{user.username}",
                    data
                )

3. 优化示例用例逻辑

ManyToManyField的add操作不会触发模型save,需手动调用:

notif = Notification.objects.create(
    content=1, 
    user_signaling=user,
    additional_info=f"{user.username} just followed you!"
)
notif.users_notified.add(other_user)
notif.save()

额外注意事项

  • 确保CHANNEL_LAYERS配置正确,生产环境建议使用Redis作为后端,避免内存后端失效。
  • 必须检查user.is_authenticated,防止匿名用户导致user.username报错。
  • 异步Consumer中查询数据库需用sync_to_async包装,否则会抛出同步/异步冲突异常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 15:56:02