如何在FastAPI中实现带实时状态提示的CSV并行下载
实现思路
HTTP请求是单向响应模式,无法在返回下载文件的同时中途更新浏览器状态,因此需要拆分流程:
- 提交查询任务到后台异步执行
- 前端通过轮询获取任务状态,根据状态显示对应提示,任务完成后触发文件下载
代码修改方案
1. 新增任务状态存储与异步任务处理
用内存字典临时存储任务状态(生产环境建议改用Redis、数据库等持久化存储),结合FastAPI的BackgroundTasks实现异步查询:
from fastapi import FastAPI, BackgroundTasks, HTTPException from fastapi.responses import StreamingResponse, JSONResponse import pandas as pd import sys import uuid app = FastAPI() # 内存存储任务状态:key为任务ID,value包含状态和数据 task_status = {} def run_query(task_id: str): try: query = "SELECT * FROM xxxx table;" cursor = conn.cursor() cursor.execute(query) df = cursor.fetch_pandas_all() cursor.close() print("Select completed") # 更新任务状态为完成并保存数据 task_status[task_id] = {"status": "completed", "data": df} except Exception as e: print("######### There was a problem in selecting data #########################") print("Error: " + str(e)) task_status[task_id] = {"status": "failed", "error": str(e)} @app.get("/submit-query") async def submit_query(background_tasks: BackgroundTasks): task_id = str(uuid.uuid4()) # 初始化任务状态为运行中 task_status[task_id] = {"status": "running"} # 添加后台执行的查询任务 background_tasks.add_task(run_query, task_id) return JSONResponse(content={"task_id": task_id})
2. 新增任务状态查询与下载接口
@app.get("/check-task-status/{task_id}") async def check_task_status(task_id: str): if task_id not in task_status: raise HTTPException(status_code=404, detail="任务不存在") status_info = task_status[task_id] if status_info["status"] == "running": return JSONResponse(content={"status": "查询运行中"}) elif status_info["status"] == "failed": return JSONResponse(content={"status": "查询失败", "error": status_info["error"]}, status_code=500) else: # 任务完成,返回状态和下载链接 return JSONResponse(content={"status": "文件下载完成", "download_url": f"/download-csv/{task_id}"}) @app.get("/download-csv/{task_id}") async def download_csv(task_id: str): if task_id not in task_status or task_status[task_id]["status"] != "completed": raise HTTPException(status_code=400, detail="任务未完成或不存在") df = task_status[task_id]["data"] # 下载完成后清理内存数据(可选) del task_status[task_id] return StreamingResponse( iter([df.to_csv(index=False)]), media_type="text/csv", headers={"Content-Disposition": f"attachment; filename=data.csv"} )
3. 前端配合逻辑(示例)
<script> async function startQuery() { // 提交查询任务 const submitRes = await fetch('/submit-query'); const { task_id } = await submitRes.json(); // 轮询任务状态 const pollInterval = setInterval(async () => { const statusRes = await fetch(`/check-task-status/${task_id}`); const statusData = await statusRes.json(); const statusElement = document.getElementById('status'); if (statusData.status === "查询运行中") { statusElement.textContent = "查询运行中"; } else if (statusData.status === "文件下载完成") { statusElement.textContent = "文件下载完成"; // 触发文件下载 window.location.href = statusData.download_url; clearInterval(pollInterval); } else if (statusData.status === "查询失败") { statusElement.textContent = `查询失败:${statusData.error}`; clearInterval(pollInterval); } }, 2000); // 每2秒轮询一次 } </script> <button onclick="startQuery()">开始下载CSV</button> <div id="status"></div>
注意事项
- 内存存储
task_status仅适合测试或小流量场景,生产环境需用Redis、数据库等持久化存储方案 - 如果查询任务耗时远超FastAPI BackgroundTasks的超时限制,建议改用Celery等专业异步任务队列
- 轮询间隔可根据实际查询耗时调整,避免过于频繁请求服务器
内容的提问来源于stack exchange,提问作者prudhvi krishna
相关产品推荐
相关产品推荐

