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

如何在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:启动服务

  1. 启动Redis服务
  2. 启动Celery Worker:
celery -A app.celery worker --loglevel=info
  1. 启动Flask应用

额外优化建议

  • 添加任务状态查询功能:用Celery任务ID,通过前端轮询或WebSocket向用户展示导入进度
  • 记录导入日志:将重复数据、错误数据存入数据库日志表,方便用户后续核对
  • 限制文件大小:在前端和后端同时设置上传文件大小上限,避免超大文件占用资源

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 12:40:54