Django单VM多用户文件处理队列解决方案咨询
解决方案:异步任务队列 + 前端轮询
核心思路
你的问题根源在于同步请求超时——浏览器和Web服务器的HTTP请求都有超时限制(通常30-60秒),不可能让用户请求挂起50秒以上等待处理完成。必须改成「任务异步处理 + 前端轮询状态」的模式,具体分三步:
- 上传接口只做任务入队,立即返回带任务ID的页面
- 用独立脚本持续监听队列,逐个处理任务
- 前端通过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>
部署说明
- 启动任务处理脚本:在VM上用
nohup python task_processor.py &让它后台运行,或者配置systemd服务确保VM重启后自动启动。 - 确保
Queue、temp、results目录有读写权限。 - 如果后续流量增大,可将JSON队列替换为Django Model(比如建一个
Task表),更可靠。
内容的提问来源于stack exchange,提问作者Syed Ali Zain Naqvi
相关产品推荐
相关产品推荐

