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

Django单VM多用户文件处理队列解决方案咨询

解决方案:异步任务队列 + 前端轮询

核心思路

你的问题根源在于同步请求超时——浏览器和Web服务器的HTTP请求都有超时限制(通常30-60秒),不可能让用户请求挂起50秒以上等待处理完成。必须改成「任务异步处理 + 前端轮询状态」的模式,具体分三步:

  1. 上传接口只做任务入队,立即返回带任务ID的页面
  2. 用独立脚本持续监听队列,逐个处理任务
  3. 前端通过Ajax轮询任务状态,完成后获取下载链接

具体实现(简单易上手版本)

1. 改造上传接口(Django视图)

上传ZIP后生成唯一任务ID,把文件存入队列目录,同时记录任务信息到队列存储(这里用JSON文件,小流量场景足够),立即返回等待页面。

import uuid
import os
import json
from django.shortcuts import render

def Uploader(request):
    if request.method == 'POST':
        file = request.FILES.get('zip_file')
        user_email = request.POST.get('user_email')
        
        # 生成唯一任务ID
        task_id = str(uuid.uuid4())
        
        # 保存上传文件到队列目录
        queue_dir = os.path.join(os.getcwd(), 'Queue')
        os.makedirs(queue_dir, exist_ok=True)
        file_path = os.path.join(queue_dir, f"{task_id}.zip")
        with open(file_path, 'wb+') as dest:
            for chunk in file.chunks():
                dest.write(chunk)
        
        # 写入队列(JSON文件存储)
        queue_file = os.path.join(os.getcwd(), 'task_queue.json')
        with open(queue_file, 'r+') as f:
            try:
                queue = json.load(f)
            except json.JSONDecodeError:
                queue = []
            queue.append({
                'task_id': task_id,
                'file_path': file_path,
                'user_email': user_email,
                'status': 'pending'
            })
            f.seek(0)
            json.dump(queue, f, indent=2)
            f.truncate()
        
        # 返回等待页面,传递任务ID
        return render(request, 'task_wait.html', {'task_id': task_id})
    
    return render(request, 'upload_form.html')

2. 编写独立任务处理脚本

这个脚本需要后台持续运行,不断检查队列,取出队首任务执行处理流程(解压、等待、生成CSV、清理),完成后标记任务状态。

import json
import time
import os
import fcntl
from shutil import unpack_archive, rmtree

def process_task(task):
    """执行你原来的workFurther逻辑"""
    # 1. 解压文件
    extract_dir = os.path.join(os.getcwd(), 'temp', task['task_id'])
    os.makedirs(extract_dir, exist_ok=True)
    unpack_archive(task['file_path'], extract_dir)
    
    # 2. 等待20秒
    time.sleep(20)
    
    # 3. 生成CSV(替换成你的实际生成逻辑)
    result_dir = os.path.join(os.getcwd(), 'results')
    os.makedirs(result_dir, exist_ok=True)
    csv_path = os.path.join(result_dir, f"{task['user_email']}.csv")
    with open(csv_path, 'w') as f:
        f.write("col1,col2\nvalue1,value2\n")  # 示例内容
    
    # 4. 清理临时文件
    rmtree(extract_dir)
    os.remove(task['file_path'])
    
    return csv_path

def main():
    queue_file = os.path.join(os.getcwd(), 'task_queue.json')
    results_file = os.path.join(os.getcwd(), 'task_results.json')
    
    while True:
        # 读取并锁定队列,避免并发处理
        task = None
        try:
            with open(queue_file, 'r+') as f:
                fcntl.flock(f, fcntl.LOCK_EX)  # Unix文件锁,Windows需替换为其他锁方式
                try:
                    queue = json.load(f)
                except json.JSONDecodeError:
                    queue = []
                
                if queue:
                    task = queue.pop(0)
                    task['status'] = 'processing'
                    # 更新队列
                    f.seek(0)
                    json.dump(queue, f, indent=2)
                    f.truncate()
                fcntl.flock(f, fcntl.LOCK_UN)
        except FileNotFoundError:
            time.sleep(5)
            continue
        
        if task:
            try:
                csv_path = process_task(task)
                # 记录任务完成状态
                with open(results_file, 'r+') as rf:
                    fcntl.flock(rf, fcntl.LOCK_EX)
                    try:
                        results = json.load(rf)
                    except json.JSONDecodeError:
                        results = {}
                    results[task['task_id']] = {
                        'status': 'completed',
                        'csv_path': csv_path,
                        'user_email': task['user_email']
                    }
                    rf.seek(0)
                    json.dump(results, rf, indent=2)
                    rf.truncate()
                    fcntl.flock(rf, fcntl.LOCK_UN)
            except Exception as e:
                # 记录任务失败状态
                with open(results_file, 'r+') as rf:
                    fcntl.flock(rf, fcntl.LOCK_EX)
                    try:
                        results = json.load(rf)
                    except json.JSONDecodeError:
                        results = {}
                    results[task['task_id']] = {
                        'status': 'failed',
                        'error': str(e)
                    }
                    rf.seek(0)
                    json.dump(results, rf, indent=2)
                    rf.truncate()
                    fcntl.flock(rf, fcntl.LOCK_UN)
        
        # 空闲时每5秒检查一次队列
        time.sleep(5)

if __name__ == '__main__':
    main()

3. 实现任务状态查询接口

供前端轮询,返回任务当前状态和下载链接(如果完成)。

import json
import os
from django.http import JsonResponse

def check_task_status(request, task_id):
    results_file = os.path.join(os.getcwd(), 'task_results.json')
    try:
        with open(results_file, 'r') as f:
            results = json.load(f)
        task_result = results.get(task_id, {'status': 'pending'})
        
        # 生成CSV下载URL(假设results目录配置为Django媒体目录)
        if task_result['status'] == 'completed':
            csv_filename = os.path.basename(task_result['csv_path'])
            task_result['csv_url'] = f"/media/results/{csv_filename}"
        
        return JsonResponse(task_result)
    except FileNotFoundError:
        return JsonResponse({'status': 'pending'})

4. 前端轮询逻辑(task_wait.html)

页面加载后每隔3秒查询一次任务状态,完成后显示下载链接。

<!DOCTYPE html>
<html>
<head>
    <title>任务处理中</title>
</head>
<body>
    <div id="status">等待任务开始...</div>
    <script>
        const taskId = "{{ task_id }}";
        const checkInterval = setInterval(() => {
            fetch(`/check-status/${taskId}/`)
                .then(res => res.json())
                .then(data => {
                    const statusEl = document.getElementById('status');
                    if (data.status === 'completed') {
                        clearInterval(checkInterval);
                        statusEl.innerHTML = `任务完成!<a href="${data.csv_url}" download>点击下载CSV</a>`;
                    } else if (data.status === 'failed') {
                        clearInterval(checkInterval);
                        statusEl.innerHTML = `任务失败:${data.error}`;
                    } else {
                        statusEl.innerHTML = `处理中,请稍候...`;
                    }
                })
                .catch(err => console.error('状态查询失败:', err));
        }, 3000);
    </script>
</body>
</html>

部署说明

  1. 启动任务处理脚本:在VM上用nohup python task_processor.py &让它后台运行,或者配置systemd服务确保VM重启后自动启动。
  2. 确保Queue、temp、results目录有读写权限。
  3. 如果后续流量增大,可将JSON队列替换为Django Model(比如建一个Task表),更可靠。

内容的提问来源于stack exchange,提问作者Syed Ali Zain Naqvi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 01:40:25