FastAPI后端向Python客户端实时推送视频处理进度实现问询
实现视频处理进度实时推送方案
核心思路
原来的上传API会阻塞到视频处理完成才返回结果,无法实现实时进度推送。我们需要把视频处理放到后台异步执行,给客户端返回唯一任务ID,再通过单独的进度查询接口让客户端获取实时进度,全程不影响视频处理流程。
一、后端代码改造(FastAPI)
1. 新增全局状态与依赖
需要用全局字典存储任务进度,同时加线程锁保证多线程操作安全,用UUID生成唯一任务ID:
import asyncio import uuid from threading import Lock, Event, Thread from fastapi import FastAPI, File, UploadFile, HTTPException, BackgroundTasks from pydantic import BaseModel import docker import os # 全局进度存储:键为任务ID,值为进度(0-100,-1表示失败) task_progress = {} # 线程锁:防止多线程同时修改进度字典导致数据异常 progress_lock = Lock() app = FastAPI() # 保留原有模型与工具函数 class VideoRequest(BaseModel): file_path: str def is_valid_video_extension(filename): # 保留你的原有实现 pass def file_exists_in_container(container, path): # 保留你的原有实现 pass # 替换你的全局变量定义 local_folder = "/path/to/local/folder" video_folder_docker = "/path/to/docker/video/folder" container_script_path = "/path/to/container/script.sh" container_name = "your-container-name"
2. 重构上传API,返回任务ID
上传完成后立即返回任务ID,把视频处理丢到后台执行:
@app.post("/api/upload_video") async def upload_video(video: UploadFile = File(...), background_tasks: BackgroundTasks = None): try: print("Received Request to upload video") if not is_valid_video_extension(video.filename): raise HTTPException(status_code=400, detail="Invalid video file extension. Supported extensions: mp4, avi, mkv, MOV") # 保存上传的视频文件 video_path = os.path.join(local_folder, video.filename) with open(video_path, 'wb') as video_file: video_file.write(video.file.read()) print("Video Uploading completed") # 生成唯一任务ID task_id = str(uuid.uuid4()) # 初始化进度为0 with progress_lock: task_progress[task_id] = 0 # 将视频处理任务加入后台队列 docker_file_path = os.path.join(video_folder_docker, video.filename) background_tasks.add_task(process_video_background, task_id, docker_file_path) return {"task_id": task_id, "message": "视频处理已启动,请用task_id查询进度"} except Exception as e: error_message = str(e) print(f"Error in upload_video: {error_message}") raise HTTPException(status_code=500, detail=f"Internal Server Error: {error_message}")
3. 后台视频处理函数,更新进度
把原来的process_video改成后台执行函数,同时更新全局进度字典:
def process_video_background(task_id: str, file_path: str): progress_exit_event = Event() try: print(f"开始处理任务:{task_id}") env_commands = ['export PATH=/demo/bins/:$PATH'] execution_command = f'{container_script_path} {file_path}' command = f"{' && '.join(env_commands)} && {execution_command}" client = docker.from_env() try: container = client.containers.get(container_name) except docker.errors.NotFound: with progress_lock: task_progress[task_id] = -1 print(f"任务{task_id}失败:Docker容器{container_name}不存在") return # 启动进度读取线程,传入任务ID更新全局进度 progress_thread = Thread(target=update_progress, args=(container, progress_exit_event, task_id)) progress_thread.start() # 执行视频处理命令 result = container.exec_run(['/bin/bash', '-c', command]) progress_exit_event.set() progress_thread.join() if result.exit_code != 0: error_message = result.output.decode('utf-8') with progress_lock: task_progress[task_id] = -1 print(f"任务{task_id}失败:{error_message}") return # 验证结果文件 json_file1_path = os.path.join(container_script_path[:-18], 'sync_events.json') json_file2_path = os.path.join(container_script_path[:-18], 'summary.json') if not file_exists_in_container(container, json_file1_path) or not file_exists_in_container(container, json_file2_path): with progress_lock: task_progress[task_id] = -1 print(f"任务{task_id}失败:未找到结果JSON文件") return # 处理完成,标记进度为100 with progress_lock: task_progress[task_id] = 100 print(f"任务{task_id}处理完成") except Exception as e: error_message = str(e) with progress_lock: task_progress[task_id] = -1 print(f"任务{task_id}处理出错:{error_message}") def update_progress(container, exit_event, task_id): try: progress_file_path = '/demo/build/progress.txt' last_progress = -2 while not exit_event.is_set(): # 读取容器内的进度文件 progress_content = container.exec_run(['/bin/cat', progress_file_path]).output.decode('utf-8').strip() try: current_progress = int(progress_content) except ValueError: # 内容不是整数,跳过 exit_event.wait(30) continue if current_progress == -1 or current_progress == last_progress: exit_event.wait(30) continue # 更新全局进度(加锁保证线程安全) with progress_lock: task_progress[task_id] = current_progress print(f"任务{task_id}进度:{current_progress}%") last_progress = current_progress # 每30秒读取一次,和文件更新频率匹配,减少Docker API调用 exit_event.wait(30) except Exception as e: print(f"任务{task_id}进度读取出错:{str(e)}")
4. 新增进度查询API
客户端通过这个接口查询指定任务的进度:
@app.get("/api/progress/{task_id}") async def get_progress(task_id: str): with progress_lock: if task_id not in task_progress: raise HTTPException(status_code=404, detail="任务不存在") progress = task_progress[task_id] if progress == -1: return {"task_id": task_id, "progress": progress, "status": "处理失败"} elif progress == 100: return {"task_id": task_id, "progress": progress, "status": "处理完成"} else: return {"task_id": task_id, "progress": progress, "status": "处理中"}
二、Python测试客户端实现(轮询方式)
前端未就绪时,用Python写一个简单客户端测试进度推送:
import requests import time # 替换为你的服务器地址 BASE_URL = "http://localhost:8000" def upload_video(file_path): print("正在上传视频...") with open(file_path, 'rb') as f: files = {"video": f} response = requests.post(f"{BASE_URL}/api/upload_video", files=files) if response.status_code == 200: return response.json()["task_id"] else: print(f"上传失败:{response.json()['detail']}") return None def check_progress(task_id): while True: response = requests.get(f"{BASE_URL}/api/progress/{task_id}") if response.status_code == 200: data = response.json() print(f"当前进度:{data['progress']}% | 状态:{data['status']}") if data['status'] in ["处理完成", "处理失败"]: break # 每30秒查询一次,和进度更新频率一致 time.sleep(30) else: print(f"查询进度失败:{response.json()['detail']}") break if __name__ == "__main__": # 替换为你的测试视频路径 test_video_path = "test_video.mp4" task_id = upload_video(test_video_path) if task_id: check_progress(task_id)
三、关键注意事项
- 线程安全:全局进度字典必须用锁保护,避免多线程同时修改导致数据异常。
- 任务清理:可以定时清理已完成/失败的任务记录(比如用
threading.Timer定期删除过期任务),避免内存占用过高。 - 性能优化:进度读取线程每30秒查询一次,和
progress.txt的更新频率匹配,减少Docker API调用次数。 - 错误标识:后台任务出错时把进度设为-1,客户端能快速识别任务失败状态。
内容的提问来源于stack exchange,提问作者Joshua
相关产品推荐
相关产品推荐

