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

RabbitMQ传输H.264视频流的发布订阅速率优化求助

问题

我正在做一个基于RabbitMQ(AMQP协议)传输H.264视频并在Web端展示的项目,流程是:捕获视频帧→编码后发送到RabbitMQ→通过Flask+Flask-SocketIO在Web端消费、解码并展示。

目前遇到RabbitMQ发布订阅性能瓶颈,消息处理速率无法超过10条/秒,达不到流畅视频流的要求,需要诊断瓶颈并给出优化方案。我已经尝试过多线程/多进程处理视频帧、优化解码函数,但问题仍未解决。

相关代码

视频捕获与发布脚本

# RabbitMQ setup
RABBITMQ_HOST = 'localhost'
EXCHANGE = 'DRONE'
CAM_LOCATION = 'Out_Front'
KEY = f'DRONE_{CAM_LOCATION}'
QUEUE_NAME = f'DRONE_{CAM_LOCATION}_video_queue'

# Path to the H.264 video file
VIDEO_FILE_PATH = 'videos/FPV.h264'

# Configure logging
logging.basicConfig(level=logging.INFO)

@contextmanager
def rabbitmq_channel(host):
    """Context manager to handle RabbitMQ channel setup and teardown."""
    connection = pika.BlockingConnection(pika.ConnectionParameters(host))
    channel = connection.channel()
    try:
        yield channel
    finally:
        connection.close()

def initialize_rabbitmq(channel):
    """Initialize RabbitMQ exchange and queue, and bind them together."""
    channel.exchange_declare(exchange=EXCHANGE, exchange_type='direct')
    channel.queue_declare(queue=QUEUE_NAME)
    channel.queue_bind(exchange=EXCHANGE, queue=QUEUE_NAME, routing_key=KEY)

def send_frame(channel, frame):
    """Encode the video frame using FFmpeg and send it to RabbitMQ."""
    ffmpeg_path = 'ffmpeg/bin/ffmpeg.exe'
    cmd = [
        ffmpeg_path,
        '-f', 'rawvideo',
        '-pix_fmt', 'rgb24',
        '-s', '{}x{}'.format(frame.shape[1], frame.shape[0]),
        '-i', 'pipe:0',
        '-f', 'h264',
        '-vcodec', 'libx264',
        '-pix_fmt', 'yuv420p',
        '-preset', 'ultrafast',
        'pipe:1'
    ]
    
    start_time = time.time()
    process = subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
    out, err = process.communicate(input=frame.tobytes())
    encoding_time = time.time() - start_time
    
    if process.returncode != 0:
        logging.error("ffmpeg error: %s", err.decode())
        raise RuntimeError("ffmpeg error")
    
    frame_size = len(out)
    logging.info("Sending frame with shape: %s, size: %d bytes", frame.shape, frame_size)
    timestamp = time.time()
    formatted_timestamp = datetime.fromtimestamp(timestamp).strftime('%H:%M:%S.%f')
    logging.info(f"Timestamp: {timestamp}") 
    logging.info(f"Formatted Timestamp: {formatted_timestamp[:-3]}")
    timestamp_bytes = struct.pack('d', timestamp)
    message_body = timestamp_bytes + out
    channel.basic_publish(exchange=EXCHANGE, routing_key=KEY, body=message_body)
    logging.info(f"Encoding time: {encoding_time:.4f} seconds")

def capture_video(channel):
    """Read video from the file, encode frames, and send them to RabbitMQ."""
    if not os.path.exists(VIDEO_FILE_PATH):
        logging.error("Error: Video file does not exist.")
        return
    cap = cv2.VideoCapture(VIDEO_FILE_PATH)
    if not cap.isOpened():
        logging.error("Error: Could not open video file.")
        return
    try:
        while True:
            start_time = time.time()
            ret, frame = cap.read()
            read_time = time.time() - start_time
            if not ret:
                break
            frame_rgb = cv2.cvtColor(frame, cv2.COLOR_BGR2RGB)
            frame_rgb = np.ascontiguousarray(frame_rgb) # Ensure the frame is contiguous
            send_frame(channel, frame_rgb)
            cv2.imshow('Video', frame)
            if cv2.waitKey(1) & 0xFF == ord('q'):
                break
            logging.info(f"Read time: {read_time:.4f} seconds")
    finally:
        cap.release()
        cv2.destroyAllWindows()

Flask后端代码

app = Flask(__name__)
CORS(app)
socketio = SocketIO(app, cors_allowed_origins="*")

RABBITMQ_HOST = 'localhost'
EXCHANGE = 'DRONE'
CAM_LOCATION = 'Out_Front'
QUEUE_NAME = f'DRONE_{CAM_LOCATION}_video_queue'

def initialize_rabbitmq():
    connection = pika.BlockingConnection(pika.ConnectionParameters(RABBITMQ_HOST))
    channel = connection.channel()
    channel.exchange_declare(exchange=EXCHANGE, exchange_type='direct')
    channel.queue_declare(queue=QUEUE_NAME)
    channel.queue_bind(exchange=EXCHANGE, queue=QUEUE_NAME, routing_key=f'DRONE_{CAM_LOCATION}')
    return connection, channel

def decode_frame(frame_data):
    # FFmpeg command to decode H.264 frame data
    ffmpeg_path = 'ffmpeg/bin/ffmpeg.exe'
    cmd = [
        ffmpeg_path,
        '-f', 'h264',
        '-i', 'pipe:0',
        '-pix_fmt', 'bgr24',
        '-vcodec', 'rawvideo',
        '-an', '-sn',
        '-f', 'rawvideo',
        'pipe:1'
    ]
    process = subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE)
    start_time = time.time()  # Start timing the decoding process
    out, err = process.communicate(input=frame_data)
    decoding_time = time.time() - start_time  # Calculate decoding time
    
    if process.returncode != 0:
        print("ffmpeg error: ", err.decode())
        return None
    frame_size = (960, 1280, 3)  # frame dimensions expected by the frontend
    frame = np.frombuffer(out, np.uint8).reshape(frame_size)
    print(f"Decoding time: {decoding_time:.4f} seconds")
    return frame

def format_timestamp(ts):
    dt = datetime.fromtimestamp(ts)
    return dt.strftime('%H:%M:%S.%f')[:-3]

def rabbitmq_consumer():
    connection, channel = initialize_rabbitmq()
    for method_frame, properties, body in channel.consume(QUEUE_NAME):
        message_receive_time = time.time()  # Time when the message is received

        # Extract the timestamp from the message body
        timestamp_bytes = body[:8]
        frame_data = body[8:]
        publish_timestamp = struct.unpack('d', timestamp_bytes)[0]

        print(f"Message Receive Time: {message_receive_time:.4f} ({format_timestamp(message_receive_time)})")
        print(f"Publish Time: {publish_timestamp:.4f} ({format_timestamp(publish_timestamp)})")

        frame = decode_frame(frame_data)
        decode_time = time.time() - message_receive_time  # Calculate decode time

        if frame is not None:
            _, buffer = cv2.imencode('.jpg', frame)
            frame_data = buffer.tobytes()
            socketio.emit('video_frame', {'frame': frame_data, 'timestamp': publish_timestamp}, namespace='/')
            emit_time = time.time()  # Time after emitting the frame

            # Log the time taken to emit the frame and its size
            rtt = emit_time - publish_timestamp  # Calculate RTT from publish to emit
            print(f"Current Time: {emit_time:.4f} ({format_timestamp(emit_time)})")
            print(f"RTT: {rtt:.4f} seconds")
            print(f"Emit time: {emit_time - message_receive_time:.4f} seconds, Frame size: {len(frame_data)} bytes")
        channel.basic_ack(method_frame.delivery_tag)

@app.route('/')
def index():
    return render_template('index.html')

@socketio.on('connect')
def handle_connect():
    print('Client connected')

@socketio.on('disconnect')
def handle_disconnect():
    print('Client disconnected')

if __name__ == '__main__':
    consumer_thread = threading.Thread(target=rabbitmq_consumer)
    consumer_thread.daemon = True
    consumer_thread.start()
    socketio.run(app, host='0.0.0.0', port=5000)

优化方案

1. 消除FFmpeg进程启动开销(核心瓶颈)

当前代码每个帧都启动一次FFmpeg进程,进程启动/销毁的开销极大,这是速率上不去的主要原因。需要改为持久化FFmpeg进程,通过管道持续传输数据:

发布端修改:

# 重构为初始化一次FFmpeg编码器,持续处理帧
def init_ffmpeg_encoder(width, height):
    ffmpeg_path = 'ffmpeg/bin/ffmpeg.exe'
    cmd = [
        ffmpeg_path,
        '-f', 'rawvideo',
        '-pix_fmt', 'rgb24',
        '-s', f'{width}x{height}',
        '-i', 'pipe:0',
        '-f', 'h264',
        '-vcodec', 'libx264',
        '-pix_fmt', 'yuv420p',
        '-preset', 'ultrafast',
        '-tune', 'zerolatency',  # 低延迟优化
        '-g', '1',  # 强制每个帧为关键帧(可选,根据需求调整)
        'pipe:1'
    ]
    return subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, bufsize=0)

def capture_video(channel):
    if not os.path.exists(VIDEO_FILE_PATH):
        logging.error("Error: Video file does not exist.")
        return
    cap = cv2.VideoCapture(VIDEO_FILE_PATH)
    if not cap.isOpened():
        logging.error("Error: Could not open video file.")
        return
    
    # 初始化一次FFmpeg编码器
    width = int(cap.get(cv2.CAP_PROP_FRAME_WIDTH))
    height = int(cap.get(cv2.CAP_PROP_FRAME_HEIGHT))
    encoder_proc = init_ffmpeg_encoder(width, height)
    
    try:
        while True:
            ret, frame = cap.read()
            if not ret:
                break
            frame_rgb = cv2.cvtColor(frame, cv2.COLOR_BGR2RGB)
            frame_rgb = np.ascontiguousarray(frame_rgb)
            
            # 直接向FFmpeg进程输入帧数据
            encoder_proc.stdin.write(frame_rgb.tobytes())
            encoder_proc.stdin.flush()
            
            # 读取编码后的H.264数据(根据实际帧大小调整读取块)
            encoded_frame = encoder_proc.stdout.read(4096)
            while len(encoded_frame) > 0:
                timestamp = time.time()
                timestamp_bytes = struct.pack('d', timestamp)
                message_body = timestamp_bytes + encoded_frame
                channel.basic_publish(exchange=EXCHANGE, routing_key=KEY, body=message_body)
                encoded_frame = encoder_proc.stdout.read(4096)
            
            cv2.imshow('Video', frame)
            if cv2.waitKey(1) & 0xFF == ord('q'):
                break
    finally:
        encoder_proc.stdin.close()
        encoder_proc.wait()
        cap.release()
        cv2.destroyAllWindows()

消费端修改:

同样持久化解码进程,避免重复启动:

def init_ffmpeg_decoder():
    ffmpeg_path = 'ffmpeg/bin/ffmpeg.exe'
    cmd = [
        ffmpeg_path,
        '-f', 'h264',
        '-i', 'pipe:0',
        '-pix_fmt', 'bgr24',
        '-vcodec', 'rawvideo',
        '-an', '-sn',
        '-f', 'rawvideo',
        'pipe:1'
    ]
    return subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, bufsize=0)

def rabbitmq_consumer():
    connection, channel = initialize_rabbitmq()
    decoder_proc = init_ffmpeg_decoder()
    frame_size = (960, 1280, 3)
    frame_byte_size = 960 * 1280 * 3
    
    # 设置预取计数,控制RabbitMQ推送消息数量
    channel.basic_qos(prefetch_count=10)
    
    try:
        for method_frame, properties, body in channel.consume(QUEUE_NAME):
            timestamp_bytes = body[:8]
            frame_data = body[8:]
            publish_timestamp = struct.unpack('d', timestamp_bytes)[0]
            
            # 向解码器输入H.264数据
            decoder_proc.stdin.write(frame_data)
            decoder_proc.stdin.flush()
            
            # 读取解码后的原始帧数据
            raw_frame = decoder_proc.stdout.read(frame_byte_size)
            if len(raw_frame) != frame_byte_size:
                print("Incomplete frame data")
                channel.basic_ack(method_frame.delivery_tag)
                continue
            
            frame = np.frombuffer(raw_frame, np.uint8).reshape(frame_size)
            # 调整JPEG质量平衡传输速度与画质
            _, buffer = cv2.imencode('.jpg', frame, [int(cv2.IMWRITE_JPEG_QUALITY), 80])
            frame_data = buffer.tobytes()
            # 用二进制模式发送SocketIO消息,减少序列化开销
            socketio.emit('video_frame', frame_data, namespace='/', binary=True)
            
            channel.basic_ack(method_frame.delivery_tag)
    finally:
        decoder_proc.stdin.close()
        decoder_proc.wait()
        connection.close()

2. 避免冗余编码(关键优化)

你的视频文件已经是H.264格式,无需用FFmpeg重新编码!直接读取文件中的H.264帧数据发送即可,完全跳过RGB转H.264的步骤:

def capture_video(channel):
    if not os.path.exists(VIDEO_FILE_PATH):
        logging.error("Error: Video file does not exist.")
        return
    
    # 按H.264 NALU分割帧(需识别起始码0x000001或0x00000001)
    with open(VIDEO_FILE_PATH, 'rb') as f:
        data = f.read()
        nalu_start = 0
        while nalu_start < len(data):
            # 查找NALU起始码
            next_start = data.find(b'\x00\x00\x01', nalu_start+3)
            if next_start == -1:
                next_start = len(data)
            frame_data = data[nalu_start:next_start]
            nalu_start = next_start
            
            timestamp = time.time()
            timestamp_bytes = struct.pack('d', timestamp)
            message_body = timestamp_bytes + frame_data
            channel.basic_publish(exchange=EXCHANGE, routing_key=KEY, body=message_body)
            
            # 控制发送速率匹配视频帧率(假设30FPS)
            time.sleep(1/30)

3. 优化RabbitMQ配置与使用

  • 批量发布消息:积累3-5帧后一次性发布,减少RabbitMQ交互次数,注意单消息大小不超过128MB的默认限制。
  • 启用批量确认:发布端开启确认模式,每发送N条消息后统一确认,减少确认开销:
    channel.confirm_delivery()
    # 每发送10条消息调用一次
    if count % 10 == 0:
        channel.wait_for_confirms()
    
  • 服务端配置优化:修改rabbitmq.conf,增加文件描述符限制、调整队列内存阈值、关闭不必要的持久化(如果允许数据丢失)。

4. 其他优化点

  • 关闭冗余日志:生产环境关闭或减少print和logging输出,避免IO开销。
  • 硬件加速编码/解码:FFmpeg添加硬件加速参数,如NVIDIA显卡使用-c:v h264_nvenc(编码)或-c:v h264_cuvid(解码)。
  • 使用异步RabbitMQ客户端:替换pika为aio-pika,提升并发处理能力。
  • 关闭UI显示:发布端的cv2.imshow会占用CPU资源,生产环境直接关闭。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 00:14:26