如何关闭Fiber端点启动的Goroutine并停止FFmpeg进程?
解决Go中FFmpeg流媒体进程的优雅停止问题
现有代码的核心问题
CombinedOutput()阻塞执行:这个方法会一直等到FFmpeg进程完全结束才返回,导致后续监听ctx取消信号的select代码完全无法在进程运行时响应取消指令。- 缺乏流状态追踪:没有保存每个流对应的
context.CancelFunc和FFmpeg进程实例,无法主动触发停止操作。 - ctx管理不当:
StreamProcess中创建子ctx时丢弃了取消函数,无法追踪和控制每个goroutine的生命周期。
解决方案步骤
1. 全局维护流状态与进程实例
定义带互斥锁的全局map,存储每个摄像头流的取消函数和FFmpeg进程,保证并发安全:
import ( "context" "os/exec" "sync" "log" "fmt" ) var ( streamMu sync.RWMutex // key: camera_id, value: 包含取消函数和进程实例的结构体 activeStreams = make(map[string]struct { cancel context.CancelFunc cmd *exec.Cmd }) )
2. 重构Stream函数,解决阻塞问题
改用Start()异步启动FFmpeg进程,同时监听ctx取消信号,确保能及时终止进程:
func Stream(meta StreamData, ctx context.Context) error { log.Printf("启动流处理: %s", meta.camera_id) // 构建FFmpeg命令 ffmpegCmd := exec.Command( "ffmpeg", "-i", meta.rtsp, "-pix_fmt", "yuv420p", "-c:v", "libx264", "-preset", "ultrafast", "-b:v", "600k", "-c:a", "aac", "-b:a", "160k", "-f", "rtsp", fmt.Sprintf("rtsp://localhost:8554/%s", meta.camera_id), ) // 启动进程(非阻塞) if err := ffmpegCmd.Start(); err != nil { log.Printf("启动FFmpeg失败 [%s]: %v", meta.camera_id, err) // 清理该流的状态 streamMu.Lock() delete(activeStreams, meta.camera_id) streamMu.Unlock() return err } // 将进程和取消函数存入全局map streamMu.Lock() activeStreams[meta.camera_id] = struct { cancel context.CancelFunc cmd *exec.Cmd }{ cancel: ctx.Value("cancel").(context.CancelFunc), cmd: ffmpegCmd, } streamMu.Unlock() // 异步等待进程结束,处理退出逻辑 go func(cameraID string) { err := ffmpegCmd.Wait() // 进程结束后清理状态 streamMu.Lock() delete(activeStreams, cameraID) streamMu.Unlock() if err != nil { log.Printf("FFmpeg进程异常退出 [%s]: %v", cameraID, err) } else { log.Printf("FFmpeg进程正常退出 [%s]", cameraID) } }(meta.camera_id) // 监听ctx取消信号,终止进程 <-ctx.Done() log.Printf("收到终止信号,杀死FFmpeg进程 [%s]", meta.camera_id) if err := ffmpegCmd.Process.Kill(); err != nil { log.Printf("杀死FFmpeg进程失败 [%s]: %v", meta.camera_id, err) } return ctx.Err() }
3. 修正StreamProcess的ctx管理
创建子ctx时保存取消函数,并传递给Stream函数:
func StreamProcess(data <-chan StreamData, parentCtx context.Context) { for v := range data { // 检查流是否已存在 streamMu.RLock() _, exists := activeStreams[v.camera_id] streamMu.RUnlock() if exists { log.Printf("流 [%s] 已存在,跳过创建", v.camera_id) continue } // 创建带取消函数的子ctx ctx, cancel := context.WithCancel(parentCtx) // 将取消函数存入ctx,方便Stream函数获取 ctx = context.WithValue(ctx, "cancel", cancel) // 启动流处理goroutine go func(meta StreamData, ctx context.Context) { defer cancel() // goroutine退出时自动取消ctx _ = Stream(meta, ctx) }(v, ctx) } }
4. 实现批量停止与单流停止函数
添加用于停止所有流或单个流的函数,直接调用即可在不关闭Fiber服务器的情况下终止FFmpeg进程:
// StopAllStreams 停止所有活跃的FFmpeg流进程 func StopAllStreams() { streamMu.Lock() defer streamMu.Unlock() for cameraID, stream := range activeStreams { log.Printf("停止流 [%s]", cameraID) stream.cancel() // 触发ctx取消 // 双重保险:直接杀死进程 if stream.cmd != nil && stream.cmd.Process != nil { if err := stream.cmd.Process.Kill(); err != nil { log.Printf("杀死进程失败 [%s]: %v", cameraID, err) } } delete(activeStreams, cameraID) } log.Println("所有流已停止") } // StopStream 停止指定摄像头的流 func StopStream(cameraID string) { streamMu.Lock() stream, exists := activeStreams[cameraID] if exists { delete(activeStreams, cameraID) } streamMu.Unlock() if !exists { log.Printf("流 [%s] 不存在", cameraID) return } log.Printf("停止流 [%s]", cameraID) stream.cancel() if stream.cmd != nil && stream.cmd.Process != nil { if err := stream.cmd.Process.Kill(); err != nil { log.Printf("杀死进程失败 [%s]: %v", cameraID, err) } } }
5. 在Fiber中调用停止函数
可以添加HTTP接口来触发停止操作,比如:
app.Get("/stop-all-streams", func(c *fiber.Ctx) error { StopAllStreams() return c.SendString("所有流已停止") }) app.Get("/stop-stream/:camera_id", func(c *fiber.Ctx) error { cameraID := c.Params("camera_id") StopStream(cameraID) return c.SendString(fmt.Sprintf("流 %s 已停止", cameraID)) })
内容的提问来源于stack exchange,提问作者Nitin
相关产品推荐
相关产品推荐

