Python asyncio调用subprocess执行ffmpeg报错‘pipe closed by peer’求助
将同步FFmpeg WebP转视频代码改为异步时的错误分析与修复
我尝试把一段能正常调用FFmpeg将目录下WebP文件转为视频的同步Python代码改成asyncio异步实现,但运行时持续报错。
可正常运行的同步代码
import subprocess import shlex from os import listdir from os.path import isfile, join output_filename = 'output.mp4' process = subprocess.Popen(shlex.split(f'ffmpeg -y -framerate 60 -i pipe: -vcodec libx265 -pix_fmt yuv420p -crf 24 output.mp4'), stdin=subprocess.PIPE) thepath = '/home/ubuntu/webpfiles/' thefiles = [f for f in listdir(thepath) if isfile(join(thepath, f))] for filename in thefiles: absolute_path = f'{thepath}{filename}' with open(absolute_path, 'rb') as f: process.stdin.write(f.read()) process.stdin.close() process.wait() process.terminate()
报错的异步代码
from os import listdir from os.path import isfile, join import shlex import asyncio outputfilename = 'output.mp4' async def write_stdin(proc): thepath = '/home/ubuntu/webpfiles/' thefiles = [f for f in listdir(thepath) if isfile(join(thepath, f))] thefiles.sort() for filename in thefiles: absolute_path = f'{thepath}{filename}' with open(absolute_path, 'rb') as f: await proc.communicate(input=f.read()) async def create_ffmpeg_subprocess(): bin = f'/home/ubuntu/bin/ffmpeg' params = f'-y -framerate 60 -i pipe: -vcodec libx265 -pix_fmt yuv420p -crf 24 {outputfilename}' proc = await asyncio.create_subprocess_exec( bin, *shlex.split(params), stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, ) return proc async def start(): loop = asyncio.get_event_loop() proc = await create_ffmpeg_subprocess() task_stdout = loop.create_task(write_stdin(proc)) await asyncio.gather(task_stdout) if __name__ == '__main__': asyncio.run(start())
错误输出
pipe closed by peer or os.write(pipe, data) raised exception. pipe closed by peer or os.write(pipe, data) raised exception. ...(每个webp文件对应一条错误信息)
错误原因分析
communicate()方法误用:proc.communicate(input=...)是一次性向进程发送所有输入,然后等待进程终止并关闭stdin/stdout/stderr。你在循环里多次调用它,第一次调用后stdin就被关闭了,后续再写入就会触发"管道关闭"错误。- 同步文件IO阻塞事件循环:
open()和f.read()是同步阻塞操作,在异步函数里调用会卡住整个事件循环,导致FFmpeg进程因长时间未收到输入而提前关闭管道。 - 未处理FFmpeg的输出缓冲区:FFmpeg的stdout/stderr如果有输出但没有被读取,缓冲区满后会阻塞进程,进而导致管道关闭。
修正后的异步代码
方案1:使用aiofiles处理异步文件读取(推荐)
先安装依赖:pip install aiofiles
from os import listdir from os.path import isfile, join import shlex import asyncio import aiofiles outputfilename = 'output.mp4' async def read_and_write(proc): thepath = '/home/ubuntu/webpfiles/' thefiles = [f for f in listdir(thepath) if isfile(join(thepath, f))] thefiles.sort() for filename in thefiles: absolute_path = f'{thepath}{filename}' # 异步读取文件 async with aiofiles.open(absolute_path, 'rb') as f: data = await f.read() # 向FFmpeg stdin写入数据并等待缓冲区排空 proc.stdin.write(data) await proc.stdin.drain() # 关闭stdin,告知FFmpeg输入结束 proc.stdin.close() await proc.stdin.wait_closed() async def consume_output(stream, name): # 持续读取进程输出,避免缓冲区阻塞 while True: data = await stream.read(4096) if not data: break # 可根据需要打印或处理输出,默认注释 # print(f"[{name}] {data.decode(errors='ignore')}") async def create_ffmpeg_subprocess(): bin = '/home/ubuntu/bin/ffmpeg' params = '-y -framerate 60 -i pipe: -vcodec libx265 -pix_fmt yuv420p -crf 24' proc = await asyncio.create_subprocess_exec( bin, *shlex.split(params), outputfilename, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, ) return proc async def start(): proc = await create_ffmpeg_subprocess() # 创建任务处理FFmpeg的stdout和stderr tasks = [ asyncio.create_task(read_and_write(proc)), asyncio.create_task(consume_output(proc.stdout, 'STDOUT')), asyncio.create_task(consume_output(proc.stderr, 'STDERR')) ] await asyncio.gather(*tasks) # 等待进程结束并获取返回码 return_code = await proc.wait() print(f"FFmpeg进程结束,返回码: {return_code}") if __name__ == '__main__': asyncio.run(start())
方案2:不依赖第三方库,用线程池处理同步文件IO
如果不想安装aiofiles,可以用asyncio.run_in_executor把同步文件读取包装成异步操作:
from os import listdir from os.path import isfile, join import shlex import asyncio from functools import partial outputfilename = 'output.mp4' def read_file_sync(absolute_path): with open(absolute_path, 'rb') as f: return f.read() async def read_and_write(proc): thepath = '/home/ubuntu/webpfiles/' thefiles = [f for f in listdir(thepath) if isfile(join(thepath, f))] thefiles.sort() loop = asyncio.get_running_loop() for filename in thefiles: absolute_path = f'{thepath}{filename}' # 用线程池执行同步文件读取 data = await loop.run_in_executor(None, partial(read_file_sync, absolute_path)) proc.stdin.write(data) await proc.stdin.drain() proc.stdin.close() await proc.stdin.wait_closed() async def consume_output(stream, name): while True: data = await stream.read(4096) if not data: break # 可根据需要打印或处理输出,默认注释 # print(f"[{name}] {data.decode(errors='ignore')}") async def create_ffmpeg_subprocess(): bin = '/home/ubuntu/bin/ffmpeg' params = '-y -framerate 60 -i pipe: -vcodec libx265 -pix_fmt yuv420p -crf 24' proc = await asyncio.create_subprocess_exec( bin, *shlex.split(params), outputfilename, stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, ) return proc async def start(): proc = await create_ffmpeg_subprocess() tasks = [ asyncio.create_task(read_and_write(proc)), asyncio.create_task(consume_output(proc.stdout, 'STDOUT')), asyncio.create_task(consume_output(proc.stderr, 'STDERR')) ] await asyncio.gather(*tasks) return_code = await proc.wait() print(f"FFmpeg进程结束,返回码: {return_code}") if __name__ == '__main__': asyncio.run(start())
关键修正点说明
- 替换
communicate()为proc.stdin.write()+await proc.stdin.drain():分批次写入数据,保持stdin打开直到所有文件发送完毕。 - 异步文件读取:避免同步IO阻塞事件循环,保证FFmpeg能持续收到输入。
- 消费FFmpeg输出:防止stdout/stderr缓冲区满导致进程阻塞,进而关闭管道。
- 正确关闭stdin并等待进程结束:确保FFmpeg处理完所有输入后正常生成视频。
内容的提问来源于stack exchange,提问作者Duke Dougal
相关产品推荐
相关产品推荐

