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()
关键修改点说明
异步子进程处理:
用asyncio.create_subprocess_exec替代subprocess.Popen,避免同步操作阻塞FastAPI事件循环,保证服务可同时处理多个WebSocket连接。确保数据已送达FFmpeg:
使用await process.stdin.drain()代替同步flush(),该方法会等待缓冲区数据完全发送到FFmpeg进程,确认当前buffer被接收后再回复客户端。修复FFmpeg关闭流程:
- 发送EOF告知FFmpeg无更多数据
- 调用
await process.wait()等待FFmpeg处理完所有数据并写入MP4元数据(moov原子),这是避免文件损坏的核心步骤 - 最后清理进程输出流资源
防止FFmpeg阻塞:
启动异步任务读取FFmpeg的stdout和stderr,避免输出缓冲区满导致FFmpeg停止处理输入数据。串行处理buffer:
循环中所有操作均使用await,确保前一个buffer的写入、处理和回复完成后,才会接收并处理下一个buffer,严格保证顺序。
额外注意事项
- 若前端发送裸视频流(如H.264),需在FFmpeg命令中添加
'-f', 'h264'指定输入格式,否则FFmpeg可能无法解析流。 - 确保
temp_dir目录存在,可添加os.makedirs(temp_dir, exist_ok=True)自动创建目录。
内容的提问来源于stack exchange,提问作者alpecca
相关产品推荐
相关产品推荐

