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

如何将Flask上传文件传递给Celery任务?序列化错误求助

解决Flask+Celery大文件上传的序列化问题

Celery默认使用JSON序列化任务参数,但FileStorage对象和gunicorn的Body对象都不属于JSON可序列化类型,直接传递会触发EncodeError,这是核心问题。以下是几种可行的解决方案:

方案1:临时文件中转(推荐,适合大文件)

思路:先在Flask端把上传的文件临时保存到服务器本地,再将临时文件路径传递给Celery任务,由Celery读取临时文件并转存到目标目录,最后可删除临时文件释放空间。这种方式避免了把大文件内容传到Celery Broker,性能更优。

修改Flask代码

import tempfile
import os
from flask import request, secure_filename

@app.route('/upload', methods=['POST'])  # 上传仅需POST方法
def upload():
    file = request.files.get("file")
    
    if not file or file.filename == "":
        return "请选择要上传的文件"
    
    filename = secure_filename(file.filename)
    
    # 创建临时文件保存上传内容
    with tempfile.NamedTemporaryFile(delete=False, suffix='.tmp') as tmp_file:
        file.save(tmp_file.name)
        tmp_file_path = tmp_file.name
    
    # 传递临时路径和目标文件名给Celery
    result = task_upload.apply_async(args=(filename, tmp_file_path), queue="upload")
    
    return result.id  # 返回任务ID供查询状态

修改Celery任务及保存方法

from pathlib import Path
import os

@celery_app.task(bind=True)
def task_upload(self, filename: str, tmp_file_path: str) -> bool:
    status = False
    try:
        status = save_file(filename, tmp_file_path)
        # 转存完成后清理临时文件
        if os.path.exists(tmp_file_path):
            os.remove(tmp_file_path)
    except Exception as e:
        print(f"Exception: {e}")
        # 出错时也尽量清理临时文件
        if os.path.exists(tmp_file_path):
            os.remove(tmp_file_path)
    return status

def save_file(filename: str, tmp_file_path: str) -> bool:
    file_path = MEDIA_DIRPATH / filename
    status = False
    
    # 分块读取临时文件并写入目标路径
    with open(tmp_file_path, 'rb') as src, open(file_path, 'wb') as dst:
        chunk_size = 4096
        while chunk := src.read(chunk_size):
            dst.write(chunk)
        status = True
    
    return status

方案2:将文件内容转为Base64字符串(不推荐大文件)

思路:把FileStorage的二进制内容转成Base64字符串(可JSON序列化),传递给Celery后再解码写入文件。但Base64会让数据体积膨胀约30%,大文件会给Celery Broker带来较大压力。

Flask端修改

import base64
from flask import request, secure_filename

@app.route('/upload', methods=['POST'])
def upload():
    file = request.files.get("file")
    
    if not file or file.filename == "":
        return "请选择要上传的文件"
    
    filename = secure_filename(file.filename)
    # 读取文件内容并转成Base64字符串
    file_content_b64 = base64.b64encode(file.read()).decode('utf-8')
    
    result = task_upload.apply_async(args=(filename, file_content_b64), queue="upload")
    
    return result.id

Celery任务及保存方法修改

import base64
from pathlib import Path

@celery_app.task(bind=True)
def task_upload(self, filename: str, file_content_b64: str) -> bool:
    status = False
    try:
        status = save_file(filename, file_content_b64)
    except Exception as e:
        print(f"Exception: {e}")
    return status

def save_file(filename: str, file_content_b64: str) -> bool:
    file_path = MEDIA_DIRPATH / filename
    status = False
    
    # 解码Base64并写入文件
    file_content = base64.b64decode(file_content_b64)
    with open(file_path, 'wb') as f:
        f.write(file_content)
    status = True
    
    return status

方案3:使用Pickle序列化(需注意安全风险)

思路:Celery支持Pickle序列化,它可以直接序列化FileStorage这类Python对象,但Pickle存在严重安全风险——如果Celery Broker被恶意访问,可能导致任意代码执行,仅适合完全信任的内部环境。

配置Celery使用Pickle

celery_app.conf.update(
    task_serializer='pickle',
    accept_content=['pickle'],
    result_serializer='pickle'
)

Flask代码(直接传递FileStorage)

from flask import request, secure_filename

@app.route('/upload', methods=['POST'])
def upload():
    file = request.files.get("file")
    
    if not file or file.filename == "":
        return "请选择要上传的文件"
    
    filename = secure_filename(file.filename)
    
    # 直接传递FileStorage对象
    result = task_upload.apply_async(args=(filename, file), queue="upload")
    
    return result.id

Celery任务及保存方法(保持原方法1逻辑)

from pathlib import Path

@celery_app.task(bind=True)
def task_upload(self, filename: str, contents) -> bool:
    status = False
    try:
        status = save_file(filename, contents)
    except Exception as e:
        print(f"Exception: {e}")
    return status

def save_file(filename: str, contents) -> bool:
    file_path = MEDIA_DIRPATH / filename
    status = False
    
    if contents:
        contents.save(file_path)
        status = True
    
    return status

优先推荐方案1,既适配大文件场景,又避免了序列化风险和Broker性能压力;小文件场景可考虑方案2;方案3仅在完全可控的内部环境中使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 12:15:33