如何将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
相关产品推荐
相关产品推荐

