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

大文件生成视频缩略图遇io: read/write on closed pipe问题求助

FFmpeg管道处理大文件时偶现io: read/write on closed pipe问题

问题背景

实现了视频转码方法输出字节数组,将其传入ThumbnailWithReader方法生成视频缩略图时,遇到io: read/write on closed pipe错误:

  • 仅在处理10MB、15MB级大文件时偶现,同一文件第3或第5次迭代才触发
  • 跳过错误后,生成的缩略图文件可以正常创建、打开和使用

核心问题

  1. 能否忽略该错误?
  2. 如何修复该错误,同时保留基于pipe的逻辑以减少不必要的读取操作?

核心代码

ThumbnailWithReader方法

func (f *FFmpeg) ThumbnailWithReader(fileBytes io.Reader, duration string) (bytes []byte, err error) {
    if fileBytes == nil {
        return nil, ErrInvalidArgument
    }

    command := "ffmpeg"
    args := []string{
        "-loglevel", "fatal",
        "-y",
        "-i", "pipe:0",
        "-ss", duration,
        "-vframes", "1",
        "-vcodec", "mjpeg",
        "-movflags", "frag_keyframe+faststart",
        "-hls_list_size", "0",
        "-f", "image2",
        "pipe:1",
    }

    return f.executeCommand(command, args, fileBytes)
}

尝试修改后的executeCommand方法

func (f *FFmpeg) executeCommand(command string, args []string, fileBytes io.Reader) (bytes []byte, err error) {
    if fileBytes == nil {
        return nil, ErrInvalidArgument
    }

    // Create a CmdRunner instance for executing the FFmpeg command.
    composer := &CmdRunner{}
    composer.Command = command
    composer.Args = args

    // Initialize input and output pipes.
    writer := composer.InitStdInPipe()
    reader := composer.InitStdOutPipe()

    // Use WaitGroup to synchronize goroutines.
    wg := &sync.WaitGroup{}

    // Goroutine for reading data from the output pipe.
    wg.Add(1)
    go func() {
        defer reader.Close()
        defer wg.Done()

        // Read data from the output pipe.
        data, errR := io.ReadAll(reader)
        // Safely update the 'bytes' variable.
        f.mutex.Lock()
        bytes = data
        err = errR
        f.mutex.Unlock()
    }()

    // Goroutine for writing data to the input pipe.
    wg.Add(1)
    //go func() {
    //  defer writer.Close()
    //  defer wg.Done()
    //
    //  // Copy data from the input to the output pipe.
    //  _, errc := io.Copy(writer, fileBytes)
    //  if errc != nil {
    //      if errors.Is(io.ErrClosedPipe, errc) {
    //          fmt.Print("close pipe\n")
    //      }
    //  }
    //}()

    go func() {
        defer wg.Done()

        // Create a buffer to store the data read from fileBytes.
        buf := make([]byte, 8192) // or any appropriate buffer size

        for {
            // Read data from fileBytes into the buffer.
            n, errr := fileBytes.Read(buf)
            if errr != nil {
                if errr == io.EOF {
                    // End of file reached, close the writer.
                    _ = writer.Close()
                    return
                } else {
                    // Other error occurred, report it.
                    fmt.Println("Error reading from fileBytes:", errr)
                    return
                }
            }

            // Write the data from the buffer to the writer.
            _, errr = writer.Write(buf[:n])
            if errr != nil {
                // Error occurred while writing, report it.
                fmt.Println("Error writing to writer:", errr)
                return
            }
        }
    }()

    // Run the FFmpeg command with pipes and wait for completion.
    err = <-composer.RunWithPipe()
    wg.Wait()

    return bytes, err
}

尝试修改后的RunWithPipe方法

func (c *CmdRunner) RunWithPipe() <-chan error {
    // Create a buffered error channel to communicate errors asynchronously.
    errCh := make(chan error, 1)

    // Prepare the command to be executed, adding some common parameters like disabling stats and setting log level to 0.
    command := c.Args

    // Create a new exec.Cmd instance with the specified command and arguments.
    cmd := exec.Command(c.Command, command...)

    // Create buffers to capture the standard output and error streams of the command.
    var outb, errb bytes.Buffer
    cmd.Stdout = &outb
    cmd.Stderr = &errb

    // If an input pipe has been set, connect it to the standard input of the command.
    if c.stdInPipeReader != nil {
        cmd.Stdin = c.stdInPipeReader
    }

    // If an output pipe has been set, connect it to the standard output of the command.
    if c.stdOutPipeWriter != nil {
        cmd.Stdout = c.stdOutPipeWriter
    }

    // Start the execution of the command asynchronously.
    err := cmd.Start()

    // Start a goroutine to handle the execution and error handling.
    go func() {
        // Check if there was an error starting the command.
        if err != nil {
            errCh <- fmt.Errorf("failed at startup execution of command, command: %v, error: %v", command, err)
            close(errCh)
            return
        } else {
            //defer c.closePipes()
            // Wait for the command to finish its execution.
            err = cmd.Wait()
            // Check if there was an error during the execution.
            if err != nil {
                err = fmt.Errorf("failed finish command execution, command: %v, error: %v", command, err)
            }
            // Close pipes to avoid resource leaks.
            go c.closePipes()

            // Send the final error status to the error channel.
            errCh <- err
        }

        // If there is any content in the standard error buffer, send it as an additional error message.
        if errb.String() != "" {
            errCh <- fmt.Errorf("stderr, error: %v", err)
        }
        close(errCh)
    }()

    // Return the error channel to the caller for asynchronous error handling.
    return errCh
}

问题解答

1. 能否忽略该错误?

可以针对性忽略,但不能无脑跳过所有错误。

这个错误的本质是:FFmpeg生成缩略图只需要视频开头的一部分数据流,当它拿到足够数据后会提前关闭输入管道,而你的写入goroutine此时还在尝试往管道中写入剩余的视频数据,因此触发io: read/write on closed pipe错误。由于缩略图已经成功生成,这个错误不会影响最终结果,但需要在代码中明确捕获io.ErrClosedPipe,避免将其他真正的管道错误(比如管道创建失败、数据读取异常)一并忽略。

2. 如何修复该错误?

需要调整管道读写的同步逻辑,处理FFmpeg提前关闭管道的场景,并修复代码中的管道覆盖问题,具体修改如下:

(1)修复写入goroutine的错误处理

将写入逻辑改回io.Copy(更简洁且高效),并针对性捕获io.ErrClosedPipe:

// 替换原写入goroutine代码
wg.Add(1)
go func() {
    defer writer.Close()
    defer wg.Done()

    _, errc := io.Copy(writer, fileBytes)
    // 仅忽略FFmpeg提前关闭管道的错误,其他错误正常返回
    if errc != nil && !errors.Is(errc, io.ErrClosedPipe) {
        f.mutex.Lock()
        if err == nil {
            err = errc
        }
        f.mutex.Unlock()
    }
}()

(2)修复RunWithPipe中的管道覆盖问题

原代码中同时设置了cmd.Stdout = &outb和cmd.Stdout = c.stdOutPipeWriter,会导致输出管道被覆盖,FFmpeg的输出无法正确传入你的读取管道。修改如下:

func (c *CmdRunner) RunWithPipe() <-chan error {
    errCh := make(chan error, 2) // 调整缓冲区大小,容纳可能的两个错误

    cmd := exec.Command(c.Command, c.Args...)

    var errb bytes.Buffer
    cmd.Stderr = &errb

    // 仅在没有设置输出管道时,才用缓冲区捕获输出
    if c.stdOutPipeWriter != nil {
        cmd.Stdout = c.stdOutPipeWriter
    } else {
        var outb bytes.Buffer
        cmd.Stdout = &outb
    }

    if c.stdInPipeReader != nil {
        cmd.Stdin = c.stdInPipeReader
    }

    err := cmd.Start()
    go func() {
        defer close(errCh)
        if err != nil {
            errCh <- fmt.Errorf("start command failed: %w", err)
            return
        }

        waitErr := cmd.Wait()
        // 关闭管道要在命令结束后执行,避免提前关闭导致错误
        c.closePipes()

        if waitErr != nil {
            errCh <- fmt.Errorf("command execution failed: %w", waitErr)
        }

        if errb.Len() > 0 {
            errCh <- fmt.Errorf("ffmpeg stderr: %s", errb.String())
        }
    }()

    return errCh
}

(3)调整命令执行与goroutine的同步顺序

原代码中先取命令错误再等待goroutine,可能导致goroutine未完成就处理错误。修改为先等待所有goroutine完成,再收集命令错误:

// 替换原executeCommand中的最后几行代码
errCh := composer.RunWithPipe()
wg.Wait()

// 收集命令执行的错误
for cmdErr := range errCh {
    if err == nil {
        err = cmdErr
    }
}

// 优先保留读取缩略图时的错误
f.mutex.Lock()
defer f.mutex.Unlock()
return bytes, err

(4)确保管道关闭逻辑正确

在CmdRunner的closePipes方法中,要正确关闭管道的两端,避免资源泄漏:

func (c *CmdRunner) closePipes() {
    if c.stdInPipeWriter != nil {
        _ = c.stdInPipeWriter.Close()
    }
    if c.stdInPipeReader != nil {
        _ = c.stdInPipeReader.Close()
    }
    if c.stdOutPipeWriter != nil {
        _ = c.stdOutPipeWriter.Close()
    }
    if c.stdOutPipeReader != nil {
        _ = c.stdOutPipeReader.Close()
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 09:37:02