Django Channels中用户离线时WebSocket通知补发方案咨询
问题描述
我使用Celery和Django Channels实现服务器端通知功能:当其他用户执行特定操作时,向目标用户推送通知,前端通过WebSocket与Consumer建立连接。当前存在核心问题:事件发生时若目标用户离线(WebSocket已关闭),会错过本应接收的通知(即使通知已存储在数据库)。需要实现的逻辑是:
- 检测目标用户在线状态,在线则即时推送通知
- 离线则将通知暂存,待用户上线后自动补发
我考虑过检查用户是否在专属Channel Layer分组中,甚至想过通过循环轮询等待用户在线再发送,但不确定方案的可行性与优雅性,特此寻求解决方法。
注:用户在线时所有功能运行正常。
现有代码
consumer.py
class MyConsumer(AsyncWebsocketConsumer): async def connect(self): self.user_inbox = f'inbox_{self.scope["user"].username}' # Check if user is anonymous and close the connection when true if self.scope["user"].is_anonymous: await self.close() self.connected_user = self.scope["user"] self.room_name = self.scope["url_route"]["kwargs"]["room_name"] self.room_group_name = "chat_%s" % self.room_name # Join room group await self.channel_layer.group_add(self.room_group_name, self.channel_name) await self.accept() async def disconnect(self, close_code): # Leave room group await self.channel_layer.group_discard(self.room_group_name, self.channel_name) async def chat_message(self, event): """FIRES WHEN MESSAGE IS SENT TO LAYERS""" event_dict = { "message_header": event["message_header"], "message_body": event["message_body"], "sender": event["sender"], "subject": event["subject"], } # Send message to WebSocket await self.send(text_data=json.dumps(event_dict))
tasks.py
@shared_task def send_notification_to_ws(ws_channel, ws_event): channel_layer = get_channel_layer() async_to_sync(channel_layer.group_send)(ws_channel, ws_event)
signals.py
ws_event = { "type": "chat_message", "message_header": message_header, "message_body": message_body, "sender": sender.username, "subject": subject, } ws_channel = "chat_%s" % recipient.username send_notification_to_ws.delay(ws_channel=ws_channel, ws_event=ws_event)
解决方案
核心思路是:离线时将通知持久化到数据库,用户上线时主动补发未读通知,同时通过Channel Layer判断用户在线状态,避免无效推送。
1. 新增Notification模型(持久化离线通知)
创建数据库模型存储未发送的通知,包含接收者、通知内容、发送状态等字段:
from django.db import models from django.contrib.auth.models import User from django.db.models import JSONField class Notification(models.Model): recipient = models.ForeignKey(User, on_delete=models.CASCADE, related_name='notifications') event_data = JSONField() # 存储完整的通知事件数据 is_sent = models.BooleanField(default=False) # 标记是否已推送 created_at = models.DateTimeField(auto_now_add=True) class Meta: ordering = ['-created_at'] # 按创建时间倒序,优先补发最新通知
2. 修改Consumer:用户上线时补发未读通知
在用户WebSocket连接成功后,自动查询并发送所有未推送的通知:
class MyConsumer(AsyncWebsocketConsumer): async def connect(self): self.user_inbox = f'inbox_{self.scope["user"].username}' if self.scope["user"].is_anonymous: await self.close() self.connected_user = self.scope["user"] self.room_name = self.scope["url_route"]["kwargs"]["room_name"] self.room_group_name = "chat_%s" % self.room_name await self.channel_layer.group_add(self.room_group_name, self.channel_name) await self.accept() # 新增:用户上线后补发未读通知 await self.send_pending_notifications() async def disconnect(self, close_code): await self.channel_layer.group_discard(self.room_group_name, self.channel_name) async def chat_message(self, event): event_dict = { "message_header": event["message_header"], "message_body": event["message_body"], "sender": event["sender"], "subject": event["subject"], } await self.send(text_data=json.dumps(event_dict)) # 新增:处理未读通知补发逻辑 async def send_pending_notifications(self): from channels.db import sync_to_async from .models import Notification # 异步查询当前用户的未发送通知 pending_notifications = await sync_to_async(list)( Notification.objects.filter(recipient=self.connected_user, is_sent=False) ) for notification in pending_notifications: # 推送通知到WebSocket await self.send(text_data=json.dumps(notification.event_data)) # 标记为已发送并保存 await sync_to_async(lambda: setattr(notification, 'is_sent', True))() await sync_to_async(notification.save)()
3. 修改Celery任务:判断用户在线状态,选择推送或存储
修改任务逻辑,先通过Channel Layer检查用户专属分组是否有活跃连接,在线则直接推送,离线则存入数据库:
from channels.layers import get_channel_layer from asgiref.sync import async_to_sync from celery import shared_task from .models import Notification from django.contrib.auth.models import User @shared_task def send_notification_to_ws(recipient_username, ws_event): channel_layer = get_channel_layer() ws_channel = f"chat_{recipient_username}" # 检查用户专属分组是否有活跃连接(判断在线状态) active_channels = async_to_sync(channel_layer.group_channels)(ws_channel) if active_channels: # 用户在线,直接推送通知到WebSocket分组 async_to_sync(channel_layer.group_send)(ws_channel, ws_event) else: # 用户离线,将通知存入数据库等待补发 recipient = User.objects.get(username=recipient_username) Notification.objects.create( recipient=recipient, event_data=ws_event )
4. 修改信号触发逻辑
更新signals.py中的调用方式,传递接收者用户名而非直接传分组名:
ws_event = { "type": "chat_message", "message_header": message_header, "message_body": message_body, "sender": sender.username, "subject": subject, } # 改为传递recipient_username send_notification_to_ws.delay(recipient_username=recipient.username, ws_event=ws_event)
优化建议
- 缓存在线用户:如果Channel Layer后端是Redis,
group_channels效率较高;若使用其他后端,可维护Redis缓存(连接时添加用户,断开时移除),通过缓存快速判断在线状态,减少Channel Layer调用。 - 清理已发送通知:可通过定时任务(如Celery Beat)定期删除已发送且超过一定时间的通知,避免数据库冗余。
- 并发处理:若用户未读通知数量大,可分批补发,避免单次推送过多导致WebSocket连接压力过大。
内容的提问来源于stack exchange,提问作者Dante
相关产品推荐
相关产品推荐

