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
相关产品推荐
相关产品推荐

