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

Celery任务先于WebSocket连接执行致消息无法送达的解决方案咨询

问题描述

用户实现了PDF文本提取的异步流程:

  1. 用户上传文件到表单,服务器验证后将文件存到本地,提交Celery任务extract_text_from_pdf,随后重定向到结果页面。
  2. 结果页面建立WebSocket连接,用于接收Celery任务处理完成后的通知。
  3. Celery任务处理完成后,通过send_notification向WebSocket群组发送提取的文本。

但遇到核心问题:Celery任务执行速度(0.14s左右)远快于客户端WebSocket连接建立的速度,导致消息已发送,但客户端还未加入群组,无法接收通知。相关日志如下:

src.pdf_processing.tasks.extract_text_from_pdf[8f535575-803d-4679-8007-7a4404a372c1] 
succeeded in 0.148360932000287s: None
src.pdf_processing.tasks.send_notification[8cbf5934-47fc-4bdb-98d8-fbdec9720ec3] 
received
src.pdf_processing.tasks.send_notification[8cbf5934-47fc-4bdb-98d8-fbdec9720ec3] 
succeeded in 0.006449124999562628s: None
WebSocket HANDSHAKING /ws/download_result/8f535575-803d-4679-8007-7a4404a372c1/ 
WebSocket CONNECT /ws/download_result/8f535575-803d-4679-8007-7a4404a372c1/ 

涉及核心代码:
表单处理视图:

def form_valid(self, form):
    storage = FileSystemStorage()
    file_path = storage.save(form.cleaned_data['file'].name, form.cleaned_data['file'])
    task = tasks.extract_text_from_pdf.delay(services.full_path(file_path))
    return HttpResponseRedirect(reverse('pdf:download_result', kwargs={'task_id': task.id}))

Celery任务:

@shared_task
def send_notification(message, task_id):
    channel_layer = get_channel_layer()

    room_group_name = f'download_result_{task_id}'

    async_to_sync(channel_layer.group_send)(
        room_group_name,
        {
            'type': 'notification_message',
            'message': message
        }
    )


@shared_task
def extract_text_from_pdf(pdf_file_path):
    with open(pdf_file_path, 'rb') as file:
        reader = PdfReader(BytesIO(file.read()))

    text = ""
    for page in reader.pages:
        text += page.extract_text() + "\n"

    text = text.replace('\n', '\r\n')

    send_notification.delay(text, current_task.request.id)
解决方案

无需更换WebSocket连接方式,通过以下几种方案即可解决消息丢失问题:

方案1:客户端连接后主动拉取结果(最可靠)

修改结果页面逻辑,先主动查询任务状态,再决定是否等待WebSocket通知:

  1. 页面加载时,通过API请求查询对应task_id的任务状态和结果。
  2. 如果任务已完成,直接展示结果;如果未完成,再建立WebSocket连接等待通知。

后端新增API视图

from django.http import JsonResponse
from celery.result import AsyncResult

def check_task_result(request, task_id):
    result = AsyncResult(task_id)
    if result.ready():
        return JsonResponse({
            'status': 'completed',
            'result': result.result
        })
    else:
        return JsonResponse({'status': 'pending'})

前端逻辑示例(JavaScript)

const taskId = window.location.pathname.split('/').pop();

// 先主动查询结果
fetch(`/api/check-task/${taskId}/`)
  .then(response => response.json())
  .then(data => {
    if (data.status === 'completed') {
      // 直接展示结果
      displayResult(data.result);
    } else {
      // 建立WebSocket等待通知
      const ws = new WebSocket(`ws://${window.location.host}/ws/download_result/${taskId}/`);
      ws.onmessage = function(event) {
        const message = event.data;
        displayResult(message);
        ws.close();
      };
    }
  });

function displayResult(text) {
  // 渲染结果到页面
  document.getElementById('result').textContent = text;
}

同时,修改Celery任务,确保结果存入Celery结果后端(需提前配置result_backend,比如Redis):

@shared_task
def extract_text_from_pdf(pdf_file_path):
    with open(pdf_file_path, 'rb') as file:
        reader = PdfReader(BytesIO(file.read()))

    text = ""
    for page in reader.pages:
        text += page.extract_text() + "\n"

    text = text.replace('\n', '\r\n')

    send_notification.delay(text, current_task.request.id)
    // 返回结果,自动存入result backend
    return text

方案2:任务端重试发送通知

修改send_notification任务,加入重试机制,检查客户端是否已加入群组,若未加入则延迟重试:

from celery.exceptions import Retry

@shared_task(bind=True, max_retries=5)
def send_notification(self, message, task_id):
    channel_layer = get_channel_layer()
    room_group_name = f'download_result_{task_id}'

    // 以Redis channel layer为例,检查群组是否有活跃成员
    from django_redis import get_redis_connection
    redis_conn = get_redis_connection('default')
    group_key = f'asgi:group:{room_group_name}'
    has_members = redis_conn.exists(group_key)

    if not has_members:
        // 无成员,延迟1秒重试
        raise self.retry(exc=Exception("No clients in group"), countdown=1)
    
    async_to_sync(channel_layer.group_send)(
        room_group_name,
        {
            'type': 'notification_message',
            'message': message
        }
    )

注意:该方案依赖channel layer的底层存储实现,仅适用于Redis等可查询的后端,不同channel layer的检查逻辑可能不同。

方案3:延迟任务执行(临时应急方案)

在extract_text_from_pdf任务中加入短暂延迟,给客户端足够时间建立WebSocket连接。但这种方式不可靠,因为网络延迟不确定:

import time

@shared_task
def extract_text_from_pdf(pdf_file_path):
    with open(pdf_file_path, 'rb') as file:
        reader = PdfReader(BytesIO(file.read()))

    text = ""
    for page in reader.pages:
        text += page.extract_text() + "\n"

    text = text.replace('\n', '\r\n')

    // 延迟1秒发送通知
    time.sleep(1)
    send_notification.delay(text, current_task.request.id)
总结

优先选择方案1,主动拉取+WebSocket通知的组合,既解决了消息丢失问题,也提升了用户体验(任务完成时无需等待WebSocket连接)。方案2适合无法修改前端逻辑的场景,方案3仅作为临时应急手段。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 11:53:16