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

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文件对应一条错误信息)

错误原因分析

  1. communicate()方法误用:proc.communicate(input=...)是一次性向进程发送所有输入,然后等待进程终止并关闭stdin/stdout/stderr。你在循环里多次调用它,第一次调用后stdin就被关闭了,后续再写入就会触发"管道关闭"错误。
  2. 同步文件IO阻塞事件循环:open()和f.read()是同步阻塞操作,在异步函数里调用会卡住整个事件循环,导致FFmpeg进程因长时间未收到输入而提前关闭管道。
  3. 未处理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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 14:15:40