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

如何在Django中使用StreamingHttpResponse批量加载聊天消息?

解决方案:分接口实现批量加载历史消息+实时流推送

要实现初始加载最新20条消息、点击加载更多批量获取更早消息,同时保留实时聊天推送的功能,最清晰高效的方式是拆分两个接口:一个负责分页加载历史消息,另一个专注于实时推送新消息。

1. 分页加载历史消息接口

这个接口支持游标分页(基于消息ID),避免偏移分页在大数据量下的性能问题,同时返回是否还有更多消息的标识。

from django.http import JsonResponse
from django.db.models import Q, F, Value, Concat
from django.utils.html import escape
import json
from django.core.serializers.json import DjangoJSONEncoder
from asgiref.sync import sync_to_async
from django.shortcuts import get_object_or_404

async def get_chat_history(request, recipient_id: int):
    recipient = await sync_to_async(get_object_or_404)(User, id=recipient_id)
    user = request.user
    before_id = request.GET.get('before_id')  # 加载更多时传入当前最早消息的ID

    # 基础查询条件:双方聊天记录、未隐藏的消息
    query = ChatMessage.objects.filter(
        Q(sender=user, recipient=recipient) | Q(sender=recipient, recipient=user)
    ).filter(
        (Q(sender=user) & Q(sender_hidden=False)) | (Q(recipient=user) & Q(recipient_hidden=False))
    ).annotate(
        profile_picture_url=Concat(
            Value(settings.MEDIA_URL),
            F("sender__userprofile__profile_picture"),
            output_field=CharField(),
        ),
        is_pinned=Q(pinned_by__in=[user]),
    ).order_by("-created_at")  # 倒序取最新的消息

    # 加载更多时,只取比before_id更早的消息
    if before_id:
        query = query.filter(id__lt=before_id)

    # 取20条数据,转为列表(异步ORM不支持直接切片转列表,需用sync_to_async包装)
    messages = await sync_to_async(list)(query[:20])

    # 处理消息格式
    processed_messages = []
    for msg in messages:
        processed_msg = {
            "id": msg.id,
            "created_at": msg.created_at.isoformat(),
            "content": escape(msg.content),
            "profile_picture_url": msg.profile_picture_url,
            "sender__id": msg.sender.id,
            "edited": msg.edited,
            "file": msg.file.url if msg.file else None,
            "file_size": msg.file_size,
            "is_pinned": msg.is_pinned,
        }
        processed_messages.append(processed_msg)

    # 判断是否还有更多消息
    has_more = len(processed_messages) == 20
    earliest_id = processed_messages[-1]["id"] if processed_messages else None

    return JsonResponse({
        "messages": processed_messages,
        "has_more": has_more,
        "earliest_id": earliest_id
    }, encoder=DjangoJSONEncoder)

2. 实时推送新消息的SSE接口

修改原有的流式接口,去掉初始加载所有历史消息的逻辑,只专注于监听并推送新产生的消息:

async def stream_chat_messages(request, recipient_id: int) -> StreamingHttpResponse:
    """仅推送用户与收件人之间的新聊天消息"""
    recipient = await sync_to_async(get_object_or_404)(User, id=recipient_id)
    user = request.user

    # 获取当前最新消息的ID,作为监听新消息的起点
    async def get_last_message_id() -> int:
        last_message = await ChatMessage.objects.filter(
            Q(sender=user, recipient=recipient) | Q(sender=recipient, recipient=user)
        ).alast()
        return last_message.id if last_message else 0

    last_id = await get_last_message_id()

    async def event_stream():
        while True:
            # 查询所有比last_id新的消息
            new_messages = (
                ChatMessage.objects.filter(
                    Q(sender=user, recipient=recipient) | Q(sender=recipient, recipient=user),
                    id__gt=last_id,
                )
                .annotate(
                    profile_picture_url=Concat(
                        Value(settings.MEDIA_URL),
                        F("sender__userprofile__profile_picture"),
                        output_field=CharField(),
                    ),
                    is_pinned=Q(pinned_by__in=[user]),
                )
                .order_by("created_at")
                .values(
                    "id", "created_at", "content", "profile_picture_url",
                    "sender__id", "edited", "file", "file_size", "is_pinned"
                )
            )

            async for msg in new_messages:
                msg["created_at"] = msg["created_at"].isoformat()
                msg["content"] = escape(msg["content"])
                if msg["file"]:
                    msg["file"] = request.build_absolute_uri(f"{settings.MEDIA_URL}{msg['file']}")
                
                json_msg = json.dumps(msg, cls=DjangoJSONEncoder)
                yield f"data: {json_msg}\n\n"
                
                nonlocal last_id
                last_id = msg["id"]
            
            await asyncio.sleep(0.1)  # 短轮询间隔,避免频繁查询数据库

    return StreamingHttpResponse(event_stream(), content_type="text/event-stream")

3. 前端配合逻辑

  • 初始加载:页面渲染完成后,调用get_chat_history接口(不带before_id),获取最新20条消息并渲染到聊天区域。
  • 实时监听:连接stream_chat_messages的SSE接口,收到新消息后添加到聊天区域底部。
  • 加载更多:点击按钮时,传入当前聊天记录中最早消息的ID作为before_id调用get_chat_history,将返回的消息插入到聊天区域顶部;如果返回的has_more为false,则隐藏“加载更多”按钮。

性能优化建议

  • 数据库索引:给ChatMessage模型添加组合索引,比如(sender, recipient, created_at)和(recipient, sender, created_at),大幅提升聊天记录查询速度。
  • 异步ORM使用:尽量用Django的异步ORM方法(如afilter、alast),减少sync_to_async的使用,提升接口响应效率。
  • 前端缓存:对已加载的历史消息做前端缓存,避免重复请求相同数据。
  • SSE重连:前端处理SSE连接断开的情况,实现自动重连逻辑,保证实时消息不丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 19:08:09