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

FastAPI+WebSocket流媒体:FFmpeg写入完成后再回复客户端的问题

解决方案

要解决你的问题,核心是确保FFmpeg完全处理每个buffer后再回复客户端,同时保证串行处理避免乱序,还要优化FFmpeg的关闭流程防止文件损坏。以下是修改后的代码和关键说明:

修改后的完整代码

import asyncio
import os
from fastapi import WebSocket, Query, Depends, APIRouter
from fastapi_jwt_auth import AuthJWT

router = APIRouter()
temp_dir = "./temp"  # 替换为你的实际临时目录

async def read_ffmpeg_stream(stream, stream_name):
    """异步读取FFmpeg的输出流,防止缓冲区满导致进程阻塞"""
    while True:
        line = await stream.readline()
        if not line:
            break
        print(f"FFmpeg {stream_name}: {line.decode().strip()}")

@router.websocket("/stream")
async def websocket_endpoint(websocket: WebSocket, token: str = Query(...), videoId: str = Query(...), authorize: AuthJWT = Depends()):
    await manager.connect(websocket)
    dataNumber = 1
    recordingFile = os.path.join(temp_dir, f"recording_{videoId}.mp4")
    
    # FFmpeg命令:若前端发送裸流(如H.264),需添加'-f h264'指定输入格式
    command = [
        'ffmpeg', 
        '-y',
        '-i', '-', 
        '-codec:v', 'copy', 
        '-f', 'mp4',
        recordingFile
    ]

    # 创建异步子进程,避免阻塞FastAPI事件循环
    process = await asyncio.create_subprocess_exec(
        *command,
        stdin=asyncio.subprocess.PIPE,
        stdout=asyncio.subprocess.PIPE,
        stderr=asyncio.subprocess.PIPE
    )

    # 启动异步任务读取FFmpeg输出,防止缓冲区溢出
    asyncio.create_task(read_ffmpeg_stream(process.stdout, "stdout"))
    asyncio.create_task(read_ffmpeg_stream(process.stderr, "stderr"))

    try:
        while True:
            data = await websocket.receive_bytes()
            if not data:
                break
            
            # 异步写入FFmpeg并等待数据被接收
            await process.stdin.write(data)
            await process.stdin.drain()  # 确保buffer已发送至FFmpeg
            
            # 确认FFmpeg接收后再回复客户端
            await websocket.send_json({"chunkNumber": dataNumber, "status": 200})
            dataNumber += 1
    except WebSocketDisconnect:
        print(f"Client disconnected: {websocket.client.host}")
    except Exception as e:
        print(f"Error occurred: {str(e)}")
    finally:
        manager.disconnect(websocket)
        # 正确关闭FFmpeg,保证MP4元数据写入
        if process.stdin.can_write_eof():
            process.stdin.write_eof()
            await process.stdin.drain()
        # 等待FFmpeg处理完所有数据并正常退出
        await process.wait()
        # 清理进程资源
        process.stdout.close()
        process.stderr.close()

关键修改点说明

  1. 异步子进程处理:
    用asyncio.create_subprocess_exec替代subprocess.Popen,避免同步操作阻塞FastAPI事件循环,保证服务可同时处理多个WebSocket连接。

  2. 确保数据已送达FFmpeg:
    使用await process.stdin.drain()代替同步flush(),该方法会等待缓冲区数据完全发送到FFmpeg进程,确认当前buffer被接收后再回复客户端。

  3. 修复FFmpeg关闭流程:

    • 发送EOF告知FFmpeg无更多数据
    • 调用await process.wait()等待FFmpeg处理完所有数据并写入MP4元数据(moov原子),这是避免文件损坏的核心步骤
    • 最后清理进程输出流资源
  4. 防止FFmpeg阻塞:
    启动异步任务读取FFmpeg的stdout和stderr,避免输出缓冲区满导致FFmpeg停止处理输入数据。

  5. 串行处理buffer:
    循环中所有操作均使用await,确保前一个buffer的写入、处理和回复完成后,才会接收并处理下一个buffer,严格保证顺序。

额外注意事项

  • 若前端发送裸视频流(如H.264),需在FFmpeg命令中添加'-f', 'h264'指定输入格式,否则FFmpeg可能无法解析流。
  • 确保temp_dir目录存在,可添加os.makedirs(temp_dir, exist_ok=True)自动创建目录。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 20:54:54