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

基于Django Rest Framework实现异步目录/文件夹上传与批量文件上传方案咨询

解决方案:基于Celery的Django DRF批量/目录文件上传改造

针对你提到的企业级系统批量文件上传需求,结合现有代码和Azure Blob存储环境,我推荐采用Celery异步处理+多文件表单上传的方案,解决同步上传的效率瓶颈,同时支持目录上传场景。下面是具体实现步骤:


一、核心问题分析

你之前尝试的批量序列化器方案失败,主要原因是:DRF的ListSerializer默认仅支持JSON格式的批量数据,无法处理multipart/form-data类型的文件上传请求——文件上传需要逐流处理,不能像普通字段那样批量解析。因此我们需要绕过序列化器的批量限制,直接在视图层处理多文件,并通过Celery异步执行上传逻辑。


二、实现步骤

1. 配置Celery异步任务队列

首先在Docker Compose中添加Redis(作为Celery的Broker和结果存储),然后在Django中配置Celery:

settings.py 新增Celery配置

CELERY_BROKER_URL = 'redis://redis:6379/0'  # 对应Docker中Redis服务名
CELERY_RESULT_BACKEND = 'redis://redis:6379/0'
CELERY_ACCEPT_CONTENT = ['json']
CELERY_TASK_SERIALIZER = 'json'
CELERY_RESULT_SERIALIZER = 'json'
CELERY_WORKER_CONCURRENCY = 4  # 根据服务器CPU核数调整

创建Celery任务文件 tasks.py

from celery import shared_task
from django.core.files import File
from django.contrib.auth.models import User
from .models import CompanyFileUpload
from .utils import Utils
import os
import tempfile
from azure.storage.blob import BlobClient

@shared_task(bind=True)
def process_single_file_upload(self, user_id, temp_file_path, file_name, comment=None):
    """异步处理单个文件上传到Azure Blob,并创建数据库记录"""
    try:
        # 1. 上传文件到Azure Blob
        blob_path = Utils.create_company_file_upload_path(None, file_name)
        with open(temp_file_path, 'rb') as f:
            file_url = default_storage.save(blob_path, File(f))
        
        # 2. 获取Azure Blob文件大小并格式化
        blob_client = BlobClient.from_connection_string(
            conn_str=os.environ.get('AZURE_STORAGE_CONNECTION_STRING'),
            container_name=os.environ.get('AZURE_STORAGE_CONTAINER_NAME'),
            blob_name=blob_path
        )
        blob_size = blob_client.get_blob_properties().size
        power = 2**10
        n = 0
        power_labels = {0: '', 1: 'Kilo', 2: 'Mega', 3: 'Giga', 4: 'Tera'}
        while blob_size > power:
            blob_size /= power
            n += 1
        formatted_size = f"{blob_size:.2f} {power_labels[n]}bytes"
        
        # 3. 创建数据库记录
        user = User.objects.get(id=user_id)
        file_upload = CompanyFileUpload.objects.create(
            user=user,
            name=file_name,
            data_file=file_url,
            comment=comment,
            file_size=formatted_size
        )
        
        # 4. 异步发送通知邮件
        send_upload_notification.delay(file_upload.id)
        
        # 清理临时文件
        os.remove(temp_file_path)
        return {"status": "success", "file_id": file_upload.id}
    
    except Exception as e:
        # 出错时清理临时文件
        if os.path.exists(temp_file_path):
            os.remove(temp_file_path)
        # 重试机制(最多3次)
        self.retry(exc=e, countdown=5, max_retries=3)
        return {"status": "failed", "error": str(e)}

@shared_task
def send_upload_notification(file_upload_id):
    """异步发送上传完成邮件"""
    try:
        file_upload = CompanyFileUpload.objects.get(id=file_upload_id)
        ctx = {
            "firstname": file_upload.user.firstname,
            "timestamp": file_upload.timestamp,
            "fullname_and_email": f"{file_upload.user.firstname} {file_upload.user.othernames}",
            "subject": "File Upload Completed",
            "recepient": file_upload.user.email,
            "email": file_upload.user.email,
            "url": file_upload.data_file.url,
            "comment": file_upload.comment
        }
        Utils.send_report_email(ctx)
    except CompanyFileUpload.DoesNotExist:
        pass

2. 修改Model与信号

  • 注释/删除原有的post_save同步邮件信号,改用Celery异步发送
  • 可以移除Modelsave方法中的文件大小计算逻辑(已在Celery任务中通过Azure API获取,更准确)

3. 实现批量/目录上传视图

创建支持多文件上传的API视图,同时兼容目录上传(前端通过webkitdirectory属性上传目录时,文件名将包含子路径):

from rest_framework.views import APIView
from rest_framework.response import Response
from rest_framework import status
from rest_framework.permissions import IsAuthenticated
from celery import group
from .tasks import process_single_file_upload
import tempfile

class BatchCompanyFileUploadAPI(APIView):
    permission_classes = [IsAuthenticated]
    
    def post(self, request, *args, **kwargs):
        # 1. 获取上传的多文件(支持目录上传,文件名将带路径)
        files = request.FILES.getlist('files')
        if not files:
            return Response({"error": "No files provided"}, status=status.HTTP_400_BAD_REQUEST)
        
        # 2. 获取可选的注释(支持单统一注释或每个文件单独注释)
        default_comment = request.data.get('comment', None)
        comments = request.data.getlist('comments', [default_comment]*len(files))
        
        # 3. 批量创建Celery任务
        task_group = []
        for idx, file in enumerate(files):
            # 写入临时文件避免内存占用过大
            with tempfile.NamedTemporaryFile(delete=False) as temp_file:
                for chunk in file.chunks():
                    temp_file.write(chunk)
                temp_path = temp_file.name
            
            # 添加单个文件处理任务
            task = process_single_file_upload.s(
                user_id=request.user.id,
                temp_file_path=temp_path,
                file_name=file.name,  # 目录上传时包含子路径,如"subdir/file.txt"
                comment=comments[idx] if idx < len(comments) else default_comment
            )
            task_group.append(task)
        
        # 4. 执行任务组并返回状态ID
        job = group(task_group)
        result = job.apply_async()
        return Response({
            "message": "Batch upload initiated",
            "group_task_id": result.id,
            "total_files": len(files)
        }, status=status.HTTP_202_ACCEPTED)

4. 新增任务状态查询视图

让前端可以查询批量上传的进度:

from celery.result import AsyncResult, GroupResult

class BatchUploadStatusAPI(APIView):
    permission_classes = [IsAuthenticated]
    
    def get(self, request, task_id, *args, **kwargs):
        result = AsyncResult(task_id)
        if isinstance(result, GroupResult):
            # 批量任务返回子任务状态
            child_statuses = []
            for child in result.children:
                child_statuses.append({
                    "task_id": child.id,
                    "status": child.status,
                    "result": child.get() if child.ready() else None
                })
            return Response({
                "group_status": result.status,
                "completed_count": sum(1 for c in result.children if c.successful()),
                "total_count": len(result.children),
                "tasks": child_statuses
            })
        else:
            # 单个任务返回详情
            return Response({
                "status": result.status,
                "result": result.get() if result.ready() else None
            })

5. 配置URL路由

from django.urls import path
from .views import (
    BatchCompanyFileUploadAPI,
    BatchUploadStatusAPI,
    # 保留原有CRUD视图
    CreateCompanyFileUploadAPI,
    ListCompanyFileUploadsAPI,
    ...
)

urlpatterns = [
    # 批量上传相关
    path('files/batch-upload/', BatchCompanyFileUploadAPI.as_view(), name='batch-file-upload'),
    path('files/batch-upload/status/<str:task_id>/', BatchUploadStatusAPI.as_view(), name='batch-upload-status'),
    # 原有单文件CRUD路由
    path('files/', CreateCompanyFileUploadAPI.as_view(), name='create-file'),
    path('files/list/', ListCompanyFileUploadsAPI.as_view(), name='list-files'),
    ...
]

三、目录上传支持

前端只需使用带有webkitdirectory属性的文件选择器:

<input type="file" id="file-input" webkitdirectory directory multiple>

上传时,每个文件的name属性会包含完整的相对路径(如docs/reports/Q3.pdf),后端的create_company_file_upload_path方法会自动将路径映射到Azure Blob的存储路径,保持目录结构。


四、性能优化建议

  1. Celery资源配置:根据服务器CPU核数调整CELERY_WORKER_CONCURRENCY,避免资源浪费
  2. Azure Blob优化:使用Azure SDK的分块上传API(upload_blob的chunk_size参数)处理超大文件
  3. 临时文件清理:确保Celery任务无论成功失败都清理临时文件,避免磁盘占用
  4. 任务监控:使用Flower监控Celery任务状态,便于排查问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 22:44:07