如何在Flask中实现后台运行批量创建客户的任务
Flask 实现Excel批量导入后台异步处理方案
核心思路
利用异步任务队列将数据导入逻辑从请求响应流程中剥离,用户提交文件后立即收到反馈,后台独立执行数据创建任务。推荐使用Celery+Redis组合实现异步任务,同时优化原代码的数据库操作效率。
步骤1:安装依赖
pip install celery redis pandas flask-sqlalchemy
步骤2:配置Celery与Flask应用
# app.py from flask import Flask from celery import Celery from flask_sqlalchemy import SQLAlchemy app = Flask(__name__) app.config['SECRET_KEY'] = 'your-secret-key' app.config['SQLALCHEMY_DATABASE_URI'] = 'your-db-uri' # Redis作为消息代理和结果存储 app.config['CELERY_BROKER_URL'] = 'redis://localhost:6379/0' app.config['CELERY_RESULT_BACKEND'] = 'redis://localhost:6379/0' db = SQLAlchemy(app) celery = Celery(app.name, broker=app.config['CELERY_BROKER_URL']) celery.conf.update(app.config)
步骤3:改写导入任务为异步函数
优化数据库操作(批量提交替代逐条commit),并添加异常处理:
# tasks.py from app import celery, db from models import Clients, Client # 替换为你的模型路径 from datetime import datetime @celery.task(bind=True, ignore_result=True) def create_client_task(self, data, institution_id): from models import Institution institution = Institution.query.get(institution_id) if not institution: return False batch_size = 100 # 每100条数据批量提交一次 batch = [] duplicate_count = 0 for idx, row in enumerate(data): try: # 提前检查CPF是否重复,避免异常捕获的性能损耗 if Client.query.filter_by(cpf=row[1]).first(): duplicate_count += 1 continue client_data = { "documents": False, "nb": str(row[0]).strip("[]"), "cpf": row[1], "rg": row[2], "name": row[3], "birth_date": row[4], "mother": row[5], "species": row[6], "wage": row[7].replace(".", ","), "address": row[8], "neighborhood": row[9], "cep": row[10], "city": row[11], "state": row[12], "phone": row[13], "dib": row[14], "bank": row[15], "agency": row[16], "account": row[17], "institution": institution, "status": "awaiting_inclusion", "upload_date": datetime.now(), "update_date": datetime.now(), "obs": row[18], } batch.append(Clients(**client_data)) # 达到批量阈值提交数据 if (idx + 1) % batch_size == 0: db.session.add_all(batch) db.session.commit() batch = [] except Exception as e: db.session.rollback() app.logger.error(f"第{idx+1}条数据处理失败: {str(e)}") continue # 提交剩余未批量的数据 if batch: db.session.add_all(batch) db.session.commit() # 可选:记录任务统计结果到数据库,供用户后续查看 return True
步骤4:修改上传视图函数
接收文件后解析数据,触发异步任务并立即返回响应:
# views.py from flask import request, redirect, url_for, flash from app import app from tasks import create_client_task import pandas as pd @app.route('/import/import-base', methods=['POST']) def import_base(): if 'files' not in request.files: flash('未选择文件') return redirect(url_for('import_page')) file = request.files['files'] if file.filename == '': flash('未选择文件') return redirect(url_for('import_page')) # 验证文件格式 allowed_extensions = {'xls', 'xlsx'} if file.filename.rsplit('.', 1)[1].lower() not in allowed_extensions: flash('仅支持XLS/XLSX格式文件') return redirect(url_for('import_page')) # 解析Excel数据 try: df = pd.read_excel(file) data = df.values.tolist() # 若Excel有表头,需跳过第一行:df = pd.read_excel(file, skiprows=1) except Exception as e: flash(f"文件解析失败: {str(e)}") return redirect(url_for('import_page')) # 替换为你的实际机构获取逻辑(比如从当前登录用户关联的机构获取) institution_id = 1 # 触发异步任务 create_client_task.delay(data, institution_id) flash('数据导入任务已启动,后台处理中,你可以继续使用系统') return redirect(url_for('dashboard')) # 跳转到用户主页
步骤5:启动服务
- 启动Redis服务
- 启动Celery Worker:
celery -A app.celery worker --loglevel=info
- 启动Flask应用
额外优化建议
- 添加任务状态查询功能:用Celery任务ID,通过前端轮询或WebSocket向用户展示导入进度
- 记录导入日志:将重复数据、错误数据存入数据库日志表,方便用户后续核对
- 限制文件大小:在前端和后端同时设置上传文件大小上限,避免超大文件占用资源
内容的提问来源于stack exchange,提问作者Gabriel Camargo
相关产品推荐
相关产品推荐

