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

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部署注意事项

  1. Redis部署:如果使用Redis,通过Helm或StatefulSet在GKE中部署Redis服务,确保Django和Celery Worker能访问到Redis地址
  2. Celery Worker部署:将Celery Worker单独打包成镜像,部署为GKE Deployment,与Django服务分离
  3. 权限配置:如果使用Google Cloud Pub/Sub,需为Worker服务账号配置Pub/Sub的订阅和发布权限
  4. 定时任务:若原有定时任务使用Celery Beat,可将Beat也部署为GKE Deployment,与Worker协同工作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 01:40:11