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

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)

三、关键注意事项

  1. 线程安全:全局进度字典必须用锁保护,避免多线程同时修改导致数据异常。
  2. 任务清理:可以定时清理已完成/失败的任务记录(比如用threading.Timer定期删除过期任务),避免内存占用过高。
  3. 性能优化:进度读取线程每30秒查询一次,和progress.txt的更新频率匹配,减少Docker API调用次数。
  4. 错误标识:后台任务出错时把进度设为-1,客户端能快速识别任务失败状态。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 17:43:12