如何在Python版aiortc中使用GStreamer对接rtspsrc与WebRTC?
GStreamer rtspsrc 对接 Python aiortc 实现WebRTC 最小实现样例
前置依赖安装
先装齐系统和Python层依赖,避免运行报错:
- 系统层(以Debian/Ubuntu为例):执行
sudo apt install gstreamer1.0-plugins-base gstreamer1.0-plugins-good gstreamer1.0-plugins-bad gstreamer1.0-plugins-ugly gstreamer1.0-libav python3-gi python3-gst-1.0 - Python层:执行
pip install aiortc av PyGObject
实现逻辑说明
核心链路非常直接:
- 用
rtspsrc插件拉取RTSP源流,经过解封装、解码、格式转换后输出原始BGR帧到appsink - 自定义aiortc的视频流轨道,持续从appsink读取最新帧,封装为WebRTC可传输的视频帧
- 基于aiortc原生PeerConnection实现SDP协商,完成WebRTC流传输
可直接运行的样板代码
import asyncio import gi import numpy as np from av import VideoFrame from aiortc import RTCPeerConnection, RTCSessionDescription, VideoStreamTrack # 初始化GStreamer gi.require_version('Gst', '1.0') from gi.repository import Gst Gst.init(None) class RTSPToWebRTCTrack(VideoStreamTrack): def __init__(self, rtsp_url: str): super().__init__() self.latest_frame = None # 构建GStreamer拉流管道 pipeline_str = f""" rtspsrc location={rtsp_url} latency=100 ! rtph264depay ! h264parse ! avdec_h264 ! videoconvert ! video/x-raw,format=BGR ! appsink name=sink emit-signals=True sync=False drop=True max-buffers=1 """ self.pipeline = Gst.parse_launch(pipeline_str) sink = self.pipeline.get_by_name("sink") # 注册帧回调,每次拿到新帧就更新最新帧缓存 def on_new_sample(sink): sample = sink.emit("pull-sample") buf = sample.get_buffer() caps = sample.get_caps() height = caps.get_structure(0).get_value("height") width = caps.get_structure(0).get_value("width") # 把GstBuffer转成numpy数组 arr = np.ndarray( (height, width, 3), buffer=buf.extract_dup(0, buf.get_size()), dtype=np.uint8 ) self.latest_frame = arr return Gst.FlowReturn.OK sink.connect("new-sample", on_new_sample) self.pipeline.set_state(Gst.State.PLAYING) async def recv(self): # 对齐WebRTC时间戳 pts, time_base = await self.next_timestamp() # 等待拿到有效帧 while self.latest_frame is None: await asyncio.sleep(0.01) # 把BGR格式帧转成aiortc可识别的VideoFrame frame = VideoFrame.from_ndarray(self.latest_frame, format="bgr24") frame.pts = pts frame.time_base = time_base return frame async def main(): # 替换成你自己的RTSP流地址 RTSP_URL = "rtsp://your_rtsp_stream_address" pc = RTCPeerConnection() # 添加RTSP转出来的视频轨道 video_track = RTSPToWebRTCTrack(RTSP_URL) pc.addTrack(video_track) # 本地测试用SDP协商逻辑:生成offer后打印,手动把answer粘进来即可完成协商 offer = await pc.createOffer() await pc.setLocalDescription(offer) print("===== 复制以下SDP Offer到对端 =====") print(pc.localDescription.sdp) print("===== 粘贴对端返回的SDP Answer,按回车确认 =====") answer_sdp = input() await pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer") ) # 保持连接运行 while True: await asyncio.sleep(3600) if __name__ == "__main__": asyncio.run(main())
常见调整点
- 如果你的RTSP流是H.265编码,把管道里的
rtph264depay ! h264parse ! avdec_h264替换为rtph265depay ! h265parse ! avdec_h265即可 latency参数控制RTSP流的抖动缓冲大小,数值越小延迟越低,弱网下容易卡顿,本地局域网测试可以设为50-200appsink上的sync=False和drop=True是低延迟关键配置,去掉会出现严重的延迟累积、帧堆积问题- 生产环境使用时把手动输入SDP的逻辑替换成WebSocket/HTTP等信令通道即可,不需要改动流处理部分的代码
注:代码运行时如果GStreamer报元件缺失错误,对照报错信息补装对应的GStreamer插件包即可
内容的提问来源于stack exchange,提问作者JAISON
相关产品推荐
相关产品推荐

