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

如何关闭Fiber端点启动的Goroutine并停止FFmpeg进程?

解决Go中FFmpeg流媒体进程的优雅停止问题

现有代码的核心问题

  1. CombinedOutput()阻塞执行:这个方法会一直等到FFmpeg进程完全结束才返回,导致后续监听ctx取消信号的select代码完全无法在进程运行时响应取消指令。
  2. 缺乏流状态追踪:没有保存每个流对应的context.CancelFunc和FFmpeg进程实例,无法主动触发停止操作。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 18:40:44