GKE部署Django调用ChatGPT API超时,如何实现异步消息队列方案?
Django异步处理ChatGPT长请求实现方案
问题背景
我在Django中构建了一个视图,该视图会向ChatGPT API发送POST请求,但当提示内容较长时,ChatGPT的响应耗时超过30秒,导致请求超时。我构思了一套异步处理方案:
- 客户端向服务器发送请求
- 服务器将请求写入消息队列,并返回一个消息ID给客户端
- 独立Worker监听消息队列,取出请求后调用OpenAI接口,再将响应写回队列
- 客户端使用消息ID定期轮询服务器获取结果
- 服务器未找到响应时返回“pending”,获取到响应后返回实际结果
我的应用部署在GKE上,且已有部分基于视图的定时任务,需要具体的实现建议。当前视图代码如下:
import openai from app.forms_prompt import PromptForm from app.models import ModelName from django.conf import settings from django.contrib.auth.decorators import login_required from django.http import HttpRequest from django.http import HttpResponse from django.shortcuts import get_object_or_404 from django.shortcuts import redirect from django.shortcuts import render from django.utils.translation import gettext_lazy as _ @login_required def form_prompt(request: HttpRequest, pk: int) -> HttpResponse: instance = get_object_or_404(ModelName, pk=pk) openai.api_key = settings.OPENAI_KEY form = PromptForm(request.POST or None, instance=setkeyword) # check if form data is valid if form.is_valid(): prompt = form.cleaned_data["text"] response = openai.ChatCompletion.create( model="gpt-4", messages=[ {"role": "user", "content": prompt}, ], ) instance.specific_field = response["choices"][0]["message"]["content"] form.save() return redirect("view_instance_name", instance.pk) return render(request, "view_prompt_name.html", context)
实现方案
一、选择消息队列工具
结合GKE环境,推荐两种方案:
- Redis + Celery:轻量易集成,适合已有定时任务使用Celery的场景
- Google Cloud Pub/Sub:GKE原生支持的托管服务,无需自行维护消息队列,适合云原生架构
二、Django端代码改造
1. 新增任务状态模型(可选,用于持久化任务)
from django.db import models from django.contrib.auth.models import User class OpenAITask(models.Model): STATUS_CHOICES = [ ('pending', '处理中'), ('completed', '已完成'), ('failed', '失败'), ] user = models.ForeignKey(User, on_delete=models.CASCADE) model_instance = models.ForeignKey('ModelName', on_delete=models.CASCADE) prompt = models.TextField() task_id = models.CharField(max_length=255, unique=True) status = models.CharField(max_length=20, choices=STATUS_CHOICES, default='pending') response = models.TextField(null=True, blank=True) created_at = models.DateTimeField(auto_now_add=True)
2. 修改原视图为异步任务提交
将同步调用OpenAI的逻辑改为提交任务到队列,返回任务ID供前端轮询:
import openai import uuid from app.forms_prompt import PromptForm from app.models import ModelName, OpenAITask from django.conf import settings from django.contrib.auth.decorators import login_required from django.http import JsonResponse, HttpRequest, HttpResponse from django.shortcuts import get_object_or_404, redirect, render from django.utils.translation import gettext_lazy as _ from django.urls import reverse # 导入Celery任务(后续定义) from app.tasks import process_openai_request @login_required def form_prompt(request: HttpRequest, pk: int) -> HttpResponse: instance = get_object_or_404(ModelName, pk=pk) form = PromptForm(request.POST or None, instance=instance) if form.is_valid(): prompt = form.cleaned_data["text"] # 生成唯一任务ID task_id = str(uuid.uuid4()) # 保存任务记录 OpenAITask.objects.create( user=request.user, model_instance=instance, prompt=prompt, task_id=task_id ) # 提交任务到Celery队列 process_openai_request.delay(task_id) # 区分AJAX请求和普通表单请求 if request.headers.get('X-Requested-With') == 'XMLHttpRequest': return JsonResponse({'task_id': task_id}) return redirect('task_status', task_id=task_id) return render(request, "view_prompt_name.html", {'form': form, 'instance': instance})
3. 新增任务状态查询视图
供前端轮询获取任务结果:
@login_required def task_status(request: HttpRequest, task_id: str) -> JsonResponse: try: task = OpenAITask.objects.get(task_id=task_id, user=request.user) if task.status == 'completed': # 更新ModelName实例字段 task.model_instance.specific_field = task.response task.model_instance.save() return JsonResponse({ 'status': 'completed', 'response': task.response, 'redirect_url': reverse('view_instance_name', args=[task.model_instance.pk]) }) elif task.status == 'failed': return JsonResponse({'status': 'failed', 'message': '请求处理失败'}) else: return JsonResponse({'status': 'pending'}) except OpenAITask.DoesNotExist: return JsonResponse({'status': 'error', 'message': '任务不存在'}, status=404)
三、Worker实现(以Celery为例)
1. Celery配置(settings.py)
# 使用Redis作为消息队列和结果存储 CELERY_BROKER_URL = 'redis://redis-service:6379/0' CELERY_RESULT_BACKEND = 'redis://redis-service:6379/0'
2. 定义Celery任务(tasks.py)
import openai import logging from celery import shared_task from django.conf import settings from app.models import OpenAITask openai.api_key = settings.OPENAI_KEY logger = logging.getLogger(__name__) @shared_task def process_openai_request(task_id: str): try: task = OpenAITask.objects.get(task_id=task_id) # 调用OpenAI API response = openai.ChatCompletion.create( model="gpt-4", messages=[ {"role": "user", "content": task.prompt}, ], ) # 更新任务状态和结果 task.response = response["choices"][0]["message"]["content"] task.status = 'completed' task.save() except Exception as e: task = OpenAITask.objects.get(task_id=task_id) task.status = 'failed' task.save() logger.error(f"任务{task_id}处理失败: {str(e)}")
四、前端轮询实现
在任务状态页面(如task_status.html)中添加JS轮询逻辑:
<div id="status">处理中...</div> <script> const taskId = "{{ task_id }}"; const pollInterval = 5000; // 每5秒轮询一次 function checkStatus() { fetch(`/task-status/${taskId}/`) .then(response => response.json()) .then(data => { const statusElement = document.getElementById('status'); if (data.status === 'completed') { statusElement.textContent = '处理完成!'; setTimeout(() => { window.location.href = data.redirect_url; }, 1000); } else if (data.status === 'failed') { statusElement.textContent = '处理失败,请重试'; clearInterval(pollTimer); } else { statusElement.textContent = '处理中...'; } }) .catch(error => console.error('请求失败:', error)); } // 启动轮询 const pollTimer = setInterval(checkStatus, pollInterval); // 页面加载时立即查询一次 checkStatus(); </script>
五、GKE部署注意事项
- Redis部署:如果使用Redis,通过Helm或StatefulSet在GKE中部署Redis服务,确保Django和Celery Worker能访问到Redis地址
- Celery Worker部署:将Celery Worker单独打包成镜像,部署为GKE Deployment,与Django服务分离
- 权限配置:如果使用Google Cloud Pub/Sub,需为Worker服务账号配置Pub/Sub的订阅和发布权限
- 定时任务:若原有定时任务使用Celery Beat,可将Beat也部署为GKE Deployment,与Worker协同工作
内容的提问来源于stack exchange,提问作者ladhari
相关产品推荐
相关产品推荐

