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

Python如何将两路RabbitMQ音视频字节流转为FFmpeg推流

两路RabbitMQ音视频流转FFmpeg广播实现方案

核心逻辑是用内存管道把RabbitMQ消费到的音视频字节直接喂给常驻FFmpeg进程,全程不落盘,延迟最低,不需要额外中间件就能对外提供直播流服务。

一、FFmpeg进程配置

首先根据你的音视频编码格式选对输入参数,两路独立流不能共用一个输入通道,用系统匿名管道分别传视频、音频数据即可。
如果你的RabbitMQ里传的是已经编码好的H264/H265视频、AAC音频(绝大多数音视频传输场景的默认选择),直接用流copy模式,几乎不占CPU,配置示例如下:

import os
import shlex
import subprocess
import asyncio
from functools import partial

# 创建两个匿名管道,分别对应视频、音频传输通道
video_pipe_read, video_pipe_write = os.pipe()
audio_pipe_read, audio_pipe_write = os.pipe()

# 低延迟+HTTP-FLV输出配置,FFmpeg直接监听端口对外提供服务,不需要额外搭流媒体服务
ffmpeg_command = f"""
ffmpeg
    -loglevel warning
    -fflags nobuffer -flags low_delay -analyzeduration 0 -probesize 32
    -f h264 -i pipe:{video_pipe_read}
    -f aac -i pipe:{audio_pipe_read}
    -c:v copy -c:a copy
    -f flv -listen 1 http://0.0.0.0:8080/live
""".replace("\n", " ")

ffmpeg_process = subprocess.Popen(
    shlex.split(ffmpeg_command),
    stdin=subprocess.DEVNULL,
    pass_fds=(video_pipe_read, audio_pipe_read)
)

# 封装异步写管道方法,避免阻塞asyncio事件循环
async def write_to_pipe(fd: int, data: bytes):
    loop = asyncio.get_running_loop()
    return await loop.run_in_executor(None, partial(os.write, fd, data))

注意:如果传输的是未编码的裸数据(比如YUV视频帧、PCM音频包),要把输入参数改成对应裸流格式:

  • 1080P 30帧YUV420P视频:-f rawvideo -pix_fmt yuv420p -s 1920x1080 -r 30
  • 44.1kHz 双声道16位PCM音频:-f s16le -ar 44100 -ac 2
    这种情况不能用copy编码,把-c:v copy -c:a copy改成-c:v libx264 -preset ultrafast -tune zerolatency -c:a aac -b:a 128k即可。

二、改造RabbitMQ消费逻辑

原有VideoProcess、AudioProcess的处理逻辑不需要改动,只要在处理完成后把字节写入对应管道即可,记得加基础的异常判断和进程存活检测:

vp = VideoProcess()
ap = AudioProcess()

async def receive_video_data(message: IncomingMessage):
    async with message.process():
        frame_data = await vp.process(frame=message.body)
        if not frame_data or ffmpeg_process.poll() is not None:
            # FFmpeg进程已退出,这里加进程重启逻辑:关闭旧管道、重建管道、拉起新FFmpeg进程即可
            return
        try:
            await write_to_pipe(video_pipe_write, frame_data)
        except BrokenPipeError:
            # 管道断裂说明FFmpeg异常退出,触发重启逻辑
            pass

async def receive_audio_data(message: IncomingMessage):
    async with message.process():
        packet_data = await ap.process(packet=message.body)
        if not packet_data or ffmpeg_process.poll() is not None:
            return
        try:
            await write_to_pipe(audio_pipe_write, packet_data)
        except BrokenPipeError:
            pass

三、对接和优化注意事项

  • 单路观看场景直接用上面的HTTP-FLV配置即可,HTML页面用flv.js拉http://服务IP:8080/live地址就能播放,端到端延迟可以控制在1-3s。
  • 多路并发观看场景不要直接用FFmpeg的-listen参数(该模式仅支持单客户端连接),把FFmpeg输出改成推流到本地轻量流媒体服务(比如SRS、nginx-rtmp),再由流媒体服务对外分发多份流即可。
  • RabbitMQ消费者要配置合理的prefetch_count,视频流建议设10-20,音频流设20-30,避免消息堆积占满内存,也防止写管道速度超过FFmpeg处理速度导致阻塞。
  • 音画不同步的话,可以在FFmpeg输入参数里加-use_wallclock_as_timestamps 1,让FFmpeg按数据接收时间打时间戳,对齐音视频。
  • 服务退出时记得关闭所有管道文件描述符,发送退出信号给FFmpeg进程,避免产生僵尸进程和文件描述符泄漏。

常见踩坑点

  • 不要试图用FFmpeg的单个stdin通道传两路流,stdin只有一个,两路流必须用独立管道传递文件描述符。
  • 低延迟场景必须加nobuffer、low_delay参数,FFmpeg默认会攒几百毫秒的缓冲探测流格式,会大幅拉高延迟。
  • 系统管道默认缓冲区大小是64KB,不要一次性写入过大的数据块,单帧大小超过缓冲区的话分块写入即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 01:15:39