大文件生成视频缩略图遇io: read/write on closed pipe问题求助
FFmpeg管道处理大文件时偶现
io: read/write on closed pipe问题 问题背景
实现了视频转码方法输出字节数组,将其传入ThumbnailWithReader方法生成视频缩略图时,遇到io: read/write on closed pipe错误:
- 仅在处理10MB、15MB级大文件时偶现,同一文件第3或第5次迭代才触发
- 跳过错误后,生成的缩略图文件可以正常创建、打开和使用
核心问题
- 能否忽略该错误?
- 如何修复该错误,同时保留基于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
相关产品推荐
相关产品推荐

