Celery任务先于WebSocket连接执行致消息无法送达的解决方案咨询
问题描述
用户实现了PDF文本提取的异步流程:
- 用户上传文件到表单,服务器验证后将文件存到本地,提交Celery任务
extract_text_from_pdf,随后重定向到结果页面。 - 结果页面建立WebSocket连接,用于接收Celery任务处理完成后的通知。
- 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通知:
- 页面加载时,通过API请求查询对应
task_id的任务状态和结果。 - 如果任务已完成,直接展示结果;如果未完成,再建立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
相关产品推荐
相关产品推荐

